415 lines
11 KiB
Go
415 lines
11 KiB
Go
// Package weatherapi adapts the internal weather API to forecast bundles.
|
|
package weatherapi
|
|
|
|
import (
|
|
"bytes"
|
|
"context"
|
|
"crypto/sha256"
|
|
"encoding/hex"
|
|
"encoding/json"
|
|
"fmt"
|
|
"io"
|
|
"net/http"
|
|
"net/url"
|
|
"path"
|
|
"strconv"
|
|
"strings"
|
|
"time"
|
|
|
|
"gitea.maximumdirect.net/eric/weatherreporter/internal/config"
|
|
"gitea.maximumdirect.net/eric/weatherreporter/internal/fileutil"
|
|
"gitea.maximumdirect.net/eric/weatherreporter/internal/forecast"
|
|
)
|
|
|
|
type Client struct {
|
|
baseURL *url.URL
|
|
httpClient *http.Client
|
|
units string
|
|
format string
|
|
timezone string
|
|
precision int
|
|
missingSource config.MissingSourceConfig
|
|
now func() time.Time
|
|
}
|
|
|
|
type Option func(*Client)
|
|
|
|
func WithHTTPClient(httpClient *http.Client) Option {
|
|
return func(c *Client) {
|
|
if httpClient != nil {
|
|
c.httpClient = httpClient
|
|
}
|
|
}
|
|
}
|
|
|
|
func WithClock(now func() time.Time) Option {
|
|
return func(c *Client) {
|
|
if now != nil {
|
|
c.now = now
|
|
}
|
|
}
|
|
}
|
|
|
|
func New(cfg config.Config, opts ...Option) (*Client, error) {
|
|
if strings.TrimSpace(cfg.WeatherAPI.BaseURL) == "" {
|
|
return nil, fmt.Errorf("weather_api.base_url is required")
|
|
}
|
|
baseURL, err := url.Parse(cfg.WeatherAPI.BaseURL)
|
|
if err != nil || baseURL.Scheme == "" || baseURL.Host == "" {
|
|
return nil, fmt.Errorf("weather_api.base_url must be an absolute URL")
|
|
}
|
|
|
|
timeout := cfg.WeatherAPI.Timeout
|
|
if timeout <= 0 {
|
|
timeout = 10 * time.Second
|
|
}
|
|
|
|
client := &Client{
|
|
baseURL: baseURL,
|
|
httpClient: &http.Client{Timeout: timeout},
|
|
units: cfg.WeatherAPI.Units,
|
|
format: cfg.WeatherAPI.Format,
|
|
timezone: cfg.WeatherAPI.Timezone,
|
|
precision: cfg.WeatherAPI.Precision,
|
|
missingSource: config.MissingSourceConfig{
|
|
Default: cfg.MissingSource.Default,
|
|
Sources: cfg.MissingSource.Sources,
|
|
},
|
|
now: time.Now,
|
|
}
|
|
for _, opt := range opts {
|
|
opt(client)
|
|
}
|
|
return client, nil
|
|
}
|
|
|
|
func (c *Client) FetchBundle(ctx context.Context) (*forecast.Bundle, error) {
|
|
fetchedAt := c.now()
|
|
builder := bundleBuilder{
|
|
client: c,
|
|
bundle: &forecast.Bundle{FetchedAt: fetchedAt},
|
|
fetchedAt: fetchedAt,
|
|
}
|
|
|
|
if err := builder.fetchObservation(ctx); err != nil {
|
|
return nil, err
|
|
}
|
|
if err := builder.fetchCurrent(ctx); err != nil {
|
|
return nil, err
|
|
}
|
|
if err := builder.fetchHourly(ctx); err != nil {
|
|
return nil, err
|
|
}
|
|
if err := builder.fetchNarrative(ctx); err != nil {
|
|
return nil, err
|
|
}
|
|
if err := builder.fetchAlerts(ctx); err != nil {
|
|
return nil, err
|
|
}
|
|
if err := builder.fetchDiscussion(ctx); err != nil {
|
|
return nil, err
|
|
}
|
|
if err := builder.addStub("daily", "daily forecast data is not available from the weather API yet"); err != nil {
|
|
return nil, err
|
|
}
|
|
if err := builder.addStub("weather_story", "NWS weather story is not available from the weather API yet"); err != nil {
|
|
return nil, err
|
|
}
|
|
|
|
return builder.bundle, nil
|
|
}
|
|
|
|
type bundleBuilder struct {
|
|
client *Client
|
|
bundle *forecast.Bundle
|
|
fetchedAt time.Time
|
|
}
|
|
|
|
func (b *bundleBuilder) fetchObservation(ctx context.Context) error {
|
|
raw, source, err := b.client.fetch(ctx, "observations", "/observations", queryOptions{precision: true})
|
|
if err != nil {
|
|
return err
|
|
}
|
|
if raw == nil {
|
|
return b.handleMissing(&source, "observation data is missing", false)
|
|
}
|
|
var observation forecast.Observation
|
|
if err := decodeSource(raw, &observation); err != nil {
|
|
return b.handleMalformed(&source, err, false)
|
|
}
|
|
source.IssuedAt = &observation.Timestamp
|
|
b.bundle.Observation = &observation
|
|
b.addSource(source)
|
|
return nil
|
|
}
|
|
|
|
func (b *bundleBuilder) fetchCurrent(ctx context.Context) error {
|
|
raw, source, err := b.client.fetch(ctx, "current", "/conditions/current", queryOptions{precision: true})
|
|
if err != nil {
|
|
return err
|
|
}
|
|
if raw == nil {
|
|
return b.handleMissing(&source, "current conditions data is missing", false)
|
|
}
|
|
var current forecast.Current
|
|
if err := decodeSource(raw, ¤t); err != nil {
|
|
return b.handleMalformed(&source, err, false)
|
|
}
|
|
b.bundle.Current = ¤t
|
|
b.addSource(source)
|
|
return nil
|
|
}
|
|
|
|
func (b *bundleBuilder) fetchHourly(ctx context.Context) error {
|
|
raw, source, err := b.client.fetch(ctx, "hourly", "/forecast/hourly", queryOptions{precision: true, timezone: true})
|
|
if err != nil {
|
|
return err
|
|
}
|
|
if raw == nil {
|
|
return b.handleMissing(&source, "hourly forecast data is missing", true)
|
|
}
|
|
var hourly forecast.ForecastRun
|
|
if err := decodeSource(raw, &hourly); err != nil {
|
|
return fmt.Errorf("decode hourly forecast from %s: %w", source.Endpoint, err)
|
|
}
|
|
if len(hourly.Periods) == 0 {
|
|
return fmt.Errorf("hourly forecast from %s contains no periods", source.Endpoint)
|
|
}
|
|
source.IssuedAt = &hourly.IssuedAt
|
|
source.UpdatedAt = hourly.UpdatedAt
|
|
b.bundle.Hourly = &hourly
|
|
b.addSource(source)
|
|
return nil
|
|
}
|
|
|
|
func (b *bundleBuilder) fetchNarrative(ctx context.Context) error {
|
|
raw, source, err := b.client.fetch(ctx, "narrative", "/forecast/narrative", queryOptions{precision: true, timezone: true})
|
|
if err != nil {
|
|
return err
|
|
}
|
|
if raw == nil {
|
|
return b.handleMissing(&source, "narrative forecast data is missing", false)
|
|
}
|
|
var narrative forecast.ForecastRun
|
|
if err := decodeSource(raw, &narrative); err != nil {
|
|
return b.handleMalformed(&source, err, false)
|
|
}
|
|
source.IssuedAt = &narrative.IssuedAt
|
|
source.UpdatedAt = narrative.UpdatedAt
|
|
b.bundle.Narrative = &narrative
|
|
b.addSource(source)
|
|
return nil
|
|
}
|
|
|
|
func (b *bundleBuilder) fetchAlerts(ctx context.Context) error {
|
|
raw, source, err := b.client.fetch(ctx, "alerts", "/alerts/active", queryOptions{allowNull: true})
|
|
if err != nil {
|
|
return err
|
|
}
|
|
if raw == nil {
|
|
return b.handleMissing(&source, "active alerts data is missing", false)
|
|
}
|
|
if isJSONNull(raw) {
|
|
b.bundle.Alerts = &forecast.AlertRun{Raw: append(json.RawMessage(nil), raw...)}
|
|
b.addSource(source)
|
|
return nil
|
|
}
|
|
var alerts forecast.AlertRun
|
|
if err := decodeSource(raw, &alerts); err != nil {
|
|
return b.handleMalformed(&source, err, false)
|
|
}
|
|
alerts.Raw = append(json.RawMessage(nil), raw...)
|
|
if alerts.AsOf != nil {
|
|
source.IssuedAt = alerts.AsOf
|
|
}
|
|
b.bundle.Alerts = &alerts
|
|
b.addSource(source)
|
|
return nil
|
|
}
|
|
|
|
func (b *bundleBuilder) fetchDiscussion(ctx context.Context) error {
|
|
raw, source, err := b.client.fetch(ctx, "discussion", "/discussion", queryOptions{timezone: true})
|
|
if err != nil {
|
|
return err
|
|
}
|
|
if raw == nil {
|
|
return b.handleMissing(&source, "forecast discussion data is missing", false)
|
|
}
|
|
var discussion forecast.Discussion
|
|
if err := decodeSource(raw, &discussion); err != nil {
|
|
return b.handleMalformed(&source, err, false)
|
|
}
|
|
source.IssuedAt = &discussion.IssuedAt
|
|
source.UpdatedAt = discussion.UpdatedAt
|
|
b.bundle.Discussion = &discussion
|
|
b.addSource(source)
|
|
return nil
|
|
}
|
|
|
|
func (b *bundleBuilder) addStub(sourceName string, message string) error {
|
|
source := forecast.Source{
|
|
Name: sourceName,
|
|
FetchedAt: b.fetchedAt,
|
|
Missing: true,
|
|
}
|
|
return b.applyMissingPolicy(&source, "missing_source", message)
|
|
}
|
|
|
|
func (b *bundleBuilder) handleMissing(source *forecast.Source, message string, required bool) error {
|
|
source.Missing = true
|
|
if required {
|
|
return fmt.Errorf("%s from %s is required", message, source.Endpoint)
|
|
}
|
|
return b.applyMissingPolicy(source, "missing_source", message)
|
|
}
|
|
|
|
func (b *bundleBuilder) handleMalformed(source *forecast.Source, err error, required bool) error {
|
|
if required {
|
|
return fmt.Errorf("decode %s from %s: %w", source.Name, source.Endpoint, err)
|
|
}
|
|
source.Missing = true
|
|
return b.applyMissingPolicy(source, "malformed_source", fmt.Sprintf("malformed %s data: %v", source.Name, err))
|
|
}
|
|
|
|
func (b *bundleBuilder) applyMissingPolicy(source *forecast.Source, code string, message string) error {
|
|
policy := b.client.policyFor(source.Name)
|
|
if policy == config.MissingSourceError {
|
|
return fmt.Errorf("%s: %s", source.Name, message)
|
|
}
|
|
if policy == config.MissingSourceWarn {
|
|
warning := forecast.SourceWarning{
|
|
Source: source.Name,
|
|
Code: code,
|
|
Severity: "warning",
|
|
Message: message,
|
|
Endpoint: source.Endpoint,
|
|
CompletenessImpact: "source omitted from bundle",
|
|
}
|
|
source.Warnings = append(source.Warnings, warning)
|
|
b.bundle.Warnings = append(b.bundle.Warnings, warning)
|
|
}
|
|
b.addSource(*source)
|
|
return nil
|
|
}
|
|
|
|
func (b *bundleBuilder) addSource(source forecast.Source) {
|
|
b.bundle.Sources = append(b.bundle.Sources, source)
|
|
}
|
|
|
|
func (c *Client) policyFor(source string) config.MissingSourcePolicy {
|
|
if policy, ok := c.missingSource.Sources[source]; ok {
|
|
return policy
|
|
}
|
|
return c.missingSource.Default
|
|
}
|
|
|
|
type queryOptions struct {
|
|
precision bool
|
|
timezone bool
|
|
allowNull bool
|
|
}
|
|
|
|
type envelope struct {
|
|
Data json.RawMessage `json:"data"`
|
|
}
|
|
|
|
func (c *Client) fetch(ctx context.Context, sourceName string, endpoint string, opts queryOptions) (json.RawMessage, forecast.Source, error) {
|
|
reqURL := c.endpointURL(endpoint, opts)
|
|
req, err := http.NewRequestWithContext(ctx, http.MethodGet, reqURL.String(), nil)
|
|
if err != nil {
|
|
return nil, forecast.Source{}, fmt.Errorf("create request for %s: %w", endpoint, err)
|
|
}
|
|
|
|
resp, err := c.httpClient.Do(req)
|
|
if err != nil {
|
|
return nil, forecast.Source{}, fmt.Errorf("fetch %s: %w", endpoint, err)
|
|
}
|
|
defer resp.Body.Close()
|
|
|
|
body, err := io.ReadAll(io.LimitReader(resp.Body, 10<<20))
|
|
if err != nil {
|
|
return nil, forecast.Source{}, fmt.Errorf("read %s response: %w", endpoint, err)
|
|
}
|
|
if resp.StatusCode < 200 || resp.StatusCode >= 300 {
|
|
return nil, forecast.Source{}, fmt.Errorf("fetch %s: unexpected HTTP status %d: %s", endpoint, resp.StatusCode, strings.TrimSpace(string(body)))
|
|
}
|
|
|
|
var env envelope
|
|
if err := json.Unmarshal(body, &env); err != nil {
|
|
return nil, forecast.Source{}, fmt.Errorf("decode %s envelope: %w", endpoint, err)
|
|
}
|
|
|
|
source := forecast.Source{
|
|
Name: sourceName,
|
|
Endpoint: endpoint,
|
|
Query: queryMap(reqURL.Query()),
|
|
FetchedAt: c.now(),
|
|
}
|
|
if len(env.Data) == 0 || (isJSONNull(env.Data) && !opts.allowNull) {
|
|
source.Missing = true
|
|
return nil, source, nil
|
|
}
|
|
hash, err := sourceHash(env.Data)
|
|
if err != nil {
|
|
return env.Data, source, nil
|
|
}
|
|
source.DataSHA256 = hash
|
|
return env.Data, source, nil
|
|
}
|
|
|
|
func isJSONNull(raw json.RawMessage) bool {
|
|
return bytes.Equal(bytes.TrimSpace(raw), []byte("null"))
|
|
}
|
|
|
|
func (c *Client) endpointURL(endpoint string, opts queryOptions) *url.URL {
|
|
reqURL := *c.baseURL
|
|
reqURL.Path = path.Join(c.baseURL.Path, endpoint)
|
|
query := reqURL.Query()
|
|
query.Set("format", c.format)
|
|
query.Set("units", c.units)
|
|
if opts.precision {
|
|
query.Set("precision", strconv.Itoa(c.precision))
|
|
}
|
|
if opts.timezone {
|
|
query.Set("tz", c.timezone)
|
|
}
|
|
reqURL.RawQuery = query.Encode()
|
|
return &reqURL
|
|
}
|
|
|
|
func queryMap(values url.Values) map[string]string {
|
|
if len(values) == 0 {
|
|
return nil
|
|
}
|
|
out := make(map[string]string, len(values))
|
|
for key, value := range values {
|
|
if len(value) > 0 {
|
|
out[key] = value[0]
|
|
}
|
|
}
|
|
return out
|
|
}
|
|
|
|
func decodeSource(raw json.RawMessage, target any) error {
|
|
if err := json.Unmarshal(raw, target); err != nil {
|
|
return err
|
|
}
|
|
return nil
|
|
}
|
|
|
|
func sourceHash(raw json.RawMessage) (string, error) {
|
|
var compact bytes.Buffer
|
|
if err := json.Compact(&compact, raw); err != nil {
|
|
return "", err
|
|
}
|
|
sum := sha256.Sum256(compact.Bytes())
|
|
return hex.EncodeToString(sum[:]), nil
|
|
}
|
|
|
|
func SaveBundle(path string, bundle *forecast.Bundle) error {
|
|
if err := fileutil.WriteJSONAtomic(path, bundle); err != nil {
|
|
return fmt.Errorf("save bundle: %w", err)
|
|
}
|
|
return nil
|
|
}
|