Add raw SPC convective outlook source
This commit is contained in:
411
internal/sources/spc/convective_outlook.go
Normal file
411
internal/sources/spc/convective_outlook.go
Normal file
@@ -0,0 +1,411 @@
|
||||
package spc
|
||||
|
||||
import (
|
||||
"context"
|
||||
"crypto/sha256"
|
||||
"encoding/json"
|
||||
"fmt"
|
||||
"math"
|
||||
"net/http"
|
||||
"strconv"
|
||||
"strings"
|
||||
"time"
|
||||
|
||||
"gitea.maximumdirect.net/ejr/feedkit/config"
|
||||
"gitea.maximumdirect.net/ejr/feedkit/event"
|
||||
fksources "gitea.maximumdirect.net/ejr/feedkit/sources"
|
||||
"gitea.maximumdirect.net/ejr/feedkit/transport"
|
||||
spcprovider "gitea.maximumdirect.net/ejr/weatherfeeder/internal/providers/spc"
|
||||
"gitea.maximumdirect.net/ejr/weatherfeeder/standards"
|
||||
)
|
||||
|
||||
const (
|
||||
driverConvectiveOutlook = "spc_convective_outlook"
|
||||
|
||||
acceptGeoJSON = "application/geo+json, application/json"
|
||||
acceptDiscussion = "text/html, application/xhtml+xml"
|
||||
acceptRSS = "application/rss+xml, application/xml, text/xml"
|
||||
)
|
||||
|
||||
type fetchProduct struct {
|
||||
Key string
|
||||
Day int
|
||||
OutlookType string
|
||||
URL string
|
||||
Accept string
|
||||
}
|
||||
|
||||
// ConvectiveOutlookSource polls SPC Day 1-3 convective outlook products and
|
||||
// emits one raw outlook bundle event.
|
||||
type ConvectiveOutlookSource struct {
|
||||
name string
|
||||
userAgent string
|
||||
locationID string
|
||||
locationName string
|
||||
latitude float64
|
||||
longitude float64
|
||||
client *http.Client
|
||||
bodyLimit int64
|
||||
|
||||
geoJSONProducts []fetchProduct
|
||||
discussions []fetchProduct
|
||||
rssURL string
|
||||
|
||||
lastHash [sha256.Size]byte
|
||||
hasHash bool
|
||||
}
|
||||
|
||||
func NewConvectiveOutlookSource(cfg config.SourceConfig) (*ConvectiveOutlookSource, error) {
|
||||
name := strings.TrimSpace(cfg.Name)
|
||||
if name == "" {
|
||||
return nil, fmt.Errorf("%s: name is required", driverConvectiveOutlook)
|
||||
}
|
||||
if cfg.Params == nil {
|
||||
return nil, fmt.Errorf("%s %q: params are required", driverConvectiveOutlook, name)
|
||||
}
|
||||
|
||||
userAgent, ok := cfg.ParamString("user_agent", "userAgent")
|
||||
if !ok {
|
||||
return nil, fmt.Errorf("%s %q: params.user_agent is required", driverConvectiveOutlook, name)
|
||||
}
|
||||
|
||||
latitude, err := requireFloatParam(cfg, "latitude")
|
||||
if err != nil {
|
||||
return nil, fmt.Errorf("%s %q: %w", driverConvectiveOutlook, name, err)
|
||||
}
|
||||
longitude, err := requireFloatParam(cfg, "longitude")
|
||||
if err != nil {
|
||||
return nil, fmt.Errorf("%s %q: %w", driverConvectiveOutlook, name, err)
|
||||
}
|
||||
|
||||
timeout := transport.DefaultHTTPTimeout
|
||||
if _, exists := cfg.Params["http_timeout"]; exists {
|
||||
var ok bool
|
||||
timeout, ok = cfg.ParamDuration("http_timeout")
|
||||
if !ok || timeout <= 0 {
|
||||
return nil, fmt.Errorf("source %q: params.http_timeout must be a positive duration", name)
|
||||
}
|
||||
}
|
||||
|
||||
bodyLimit := transport.DefaultHTTPResponseBodyLimitBytes
|
||||
if _, exists := cfg.Params["http_response_body_limit_bytes"]; exists {
|
||||
rawLimit, ok := cfg.ParamInt("http_response_body_limit_bytes")
|
||||
if !ok || rawLimit <= 0 {
|
||||
return nil, fmt.Errorf("source %q: params.http_response_body_limit_bytes must be a positive integer", name)
|
||||
}
|
||||
bodyLimit = int64(rawLimit)
|
||||
}
|
||||
|
||||
geoJSONProducts, err := configuredGeoJSONProducts(cfg)
|
||||
if err != nil {
|
||||
return nil, fmt.Errorf("%s %q: %w", driverConvectiveOutlook, name, err)
|
||||
}
|
||||
discussions, err := configuredDiscussionProducts(cfg)
|
||||
if err != nil {
|
||||
return nil, fmt.Errorf("%s %q: %w", driverConvectiveOutlook, name, err)
|
||||
}
|
||||
|
||||
rssURL := ""
|
||||
if rawRSSURL, ok := cfg.ParamString("rss_url", "rssURL"); ok {
|
||||
rssURL = rawRSSURL
|
||||
}
|
||||
|
||||
locationID, _ := cfg.ParamString("location_id", "locationID")
|
||||
locationName, _ := cfg.ParamString("location_name", "locationName")
|
||||
|
||||
return &ConvectiveOutlookSource{
|
||||
name: name,
|
||||
userAgent: userAgent,
|
||||
locationID: locationID,
|
||||
locationName: locationName,
|
||||
latitude: latitude,
|
||||
longitude: longitude,
|
||||
client: transport.NewHTTPClient(timeout),
|
||||
bodyLimit: bodyLimit,
|
||||
geoJSONProducts: geoJSONProducts,
|
||||
discussions: discussions,
|
||||
rssURL: rssURL,
|
||||
}, nil
|
||||
}
|
||||
|
||||
func (s *ConvectiveOutlookSource) Name() string { return s.name }
|
||||
|
||||
func (s *ConvectiveOutlookSource) Kinds() []event.Kind {
|
||||
return []event.Kind{event.Kind("outlook")}
|
||||
}
|
||||
|
||||
func (s *ConvectiveOutlookSource) Poll(ctx context.Context) ([]event.Event, error) {
|
||||
fetchedAt := time.Now().UTC()
|
||||
bundle := spcprovider.RawConvectiveOutlookBundle{
|
||||
LocationID: s.locationID,
|
||||
LocationName: s.locationName,
|
||||
Latitude: s.latitude,
|
||||
Longitude: s.longitude,
|
||||
FetchedAt: fetchedAt,
|
||||
Products: make([]spcprovider.RawOutlookProduct, 0, len(s.geoJSONProducts)),
|
||||
Discussions: make([]spcprovider.RawDiscussionPage, 0, len(s.discussions)),
|
||||
}
|
||||
hash := sha256.New()
|
||||
|
||||
var latestIssue time.Time
|
||||
for _, product := range s.geoJSONProducts {
|
||||
body, err := s.fetch(ctx, product.URL, product.Accept)
|
||||
if err != nil {
|
||||
return nil, fmt.Errorf("fetch geojson %s: %w", product.Key, err)
|
||||
}
|
||||
addHashPart(hash, product.Key, product.URL, body)
|
||||
|
||||
bundle.Products = append(bundle.Products, spcprovider.RawOutlookProduct{
|
||||
Key: product.Key,
|
||||
Day: product.Day,
|
||||
OutlookType: product.OutlookType,
|
||||
URL: product.URL,
|
||||
FetchedAt: fetchedAt,
|
||||
Body: json.RawMessage(body),
|
||||
})
|
||||
if t := latestIssueTime(body); !t.IsZero() && (latestIssue.IsZero() || t.After(latestIssue)) {
|
||||
latestIssue = t
|
||||
}
|
||||
}
|
||||
|
||||
var latestUpdated time.Time
|
||||
for _, product := range s.discussions {
|
||||
body, err := s.fetch(ctx, product.URL, product.Accept)
|
||||
if err != nil {
|
||||
return nil, fmt.Errorf("fetch discussion %s: %w", product.Key, err)
|
||||
}
|
||||
addHashPart(hash, product.Key, product.URL, body)
|
||||
|
||||
bodyText := string(body)
|
||||
bundle.Discussions = append(bundle.Discussions, spcprovider.RawDiscussionPage{
|
||||
Key: product.Key,
|
||||
Day: product.Day,
|
||||
URL: product.URL,
|
||||
FetchedAt: fetchedAt,
|
||||
Body: bodyText,
|
||||
})
|
||||
if t := discussionUpdatedTime(bodyText); !t.IsZero() && (latestUpdated.IsZero() || t.After(latestUpdated)) {
|
||||
latestUpdated = t
|
||||
}
|
||||
}
|
||||
|
||||
var rssBuild time.Time
|
||||
if s.rssURL != "" {
|
||||
body, err := s.fetch(ctx, s.rssURL, acceptRSS)
|
||||
if err != nil {
|
||||
return nil, fmt.Errorf("fetch rss: %w", err)
|
||||
}
|
||||
addHashPart(hash, "rss", s.rssURL, body)
|
||||
|
||||
bodyText := string(body)
|
||||
bundle.RSS = &spcprovider.RawRSSFeed{
|
||||
URL: s.rssURL,
|
||||
FetchedAt: fetchedAt,
|
||||
Body: bodyText,
|
||||
}
|
||||
if feed, err := spcprovider.ParseRSSFeed(bodyText); err == nil && feed.LastBuildDate != nil {
|
||||
rssBuild = feed.LastBuildDate.UTC()
|
||||
}
|
||||
}
|
||||
|
||||
var currentHash [sha256.Size]byte
|
||||
copy(currentHash[:], hash.Sum(nil))
|
||||
if s.hasHash && currentHash == s.lastHash {
|
||||
return nil, nil
|
||||
}
|
||||
s.lastHash = currentHash
|
||||
s.hasHash = true
|
||||
|
||||
effectiveAt := chooseEffectiveTime(latestIssue, latestUpdated, rssBuild, fetchedAt)
|
||||
emittedAt := time.Now().UTC()
|
||||
eventID := fksources.DefaultEventID("", s.name, &effectiveAt, emittedAt)
|
||||
|
||||
return fksources.SingleEvent(
|
||||
event.Kind("outlook"),
|
||||
s.name,
|
||||
standards.SchemaRawSPCConvectiveOutlookV1,
|
||||
eventID,
|
||||
emittedAt,
|
||||
&effectiveAt,
|
||||
bundle,
|
||||
)
|
||||
}
|
||||
|
||||
func (s *ConvectiveOutlookSource) fetch(ctx context.Context, url, accept string) ([]byte, error) {
|
||||
return transport.FetchBodyWithLimit(ctx, s.client, url, s.userAgent, accept, s.bodyLimit)
|
||||
}
|
||||
|
||||
func configuredGeoJSONProducts(cfg config.SourceConfig) ([]fetchProduct, error) {
|
||||
overrides, err := optionalStringMap(cfg, "geojson_urls")
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
out := make([]fetchProduct, 0, len(spcprovider.GeoJSONProducts()))
|
||||
for _, product := range spcprovider.GeoJSONProducts() {
|
||||
url := product.URL
|
||||
if override := strings.TrimSpace(overrides[product.Key]); override != "" {
|
||||
url = override
|
||||
}
|
||||
out = append(out, fetchProduct{
|
||||
Key: product.Key,
|
||||
Day: product.Day,
|
||||
OutlookType: product.OutlookType,
|
||||
URL: url,
|
||||
Accept: acceptGeoJSON,
|
||||
})
|
||||
}
|
||||
return out, nil
|
||||
}
|
||||
|
||||
func configuredDiscussionProducts(cfg config.SourceConfig) ([]fetchProduct, error) {
|
||||
overrides, err := optionalStringMap(cfg, "discussion_urls")
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
out := make([]fetchProduct, 0, len(spcprovider.DiscussionProducts()))
|
||||
for _, product := range spcprovider.DiscussionProducts() {
|
||||
url := product.URL
|
||||
if override := strings.TrimSpace(overrides[product.Key]); override != "" {
|
||||
url = override
|
||||
}
|
||||
out = append(out, fetchProduct{
|
||||
Key: product.Key,
|
||||
Day: product.Day,
|
||||
URL: url,
|
||||
Accept: acceptDiscussion,
|
||||
})
|
||||
}
|
||||
return out, nil
|
||||
}
|
||||
|
||||
func optionalStringMap(cfg config.SourceConfig, key string) (map[string]string, error) {
|
||||
raw, ok := cfg.Params[key]
|
||||
if !ok || raw == nil {
|
||||
return map[string]string{}, nil
|
||||
}
|
||||
|
||||
out := map[string]string{}
|
||||
switch typed := raw.(type) {
|
||||
case map[string]string:
|
||||
for k, v := range typed {
|
||||
if strings.TrimSpace(k) != "" && strings.TrimSpace(v) != "" {
|
||||
out[strings.TrimSpace(k)] = strings.TrimSpace(v)
|
||||
}
|
||||
}
|
||||
case map[string]any:
|
||||
for k, v := range typed {
|
||||
s, ok := v.(string)
|
||||
if !ok {
|
||||
return nil, fmt.Errorf("params.%s[%q] must be a string", key, k)
|
||||
}
|
||||
if strings.TrimSpace(k) != "" && strings.TrimSpace(s) != "" {
|
||||
out[strings.TrimSpace(k)] = strings.TrimSpace(s)
|
||||
}
|
||||
}
|
||||
default:
|
||||
return nil, fmt.Errorf("params.%s must be a string map", key)
|
||||
}
|
||||
return out, nil
|
||||
}
|
||||
|
||||
func requireFloatParam(cfg config.SourceConfig, key string) (float64, error) {
|
||||
raw, ok := cfg.Params[key]
|
||||
if !ok || raw == nil {
|
||||
return 0, fmt.Errorf("params.%s is required", key)
|
||||
}
|
||||
v, ok := numberFromAny(raw)
|
||||
if !ok {
|
||||
return 0, fmt.Errorf("params.%s must be a number", key)
|
||||
}
|
||||
if math.IsNaN(v) || math.IsInf(v, 0) {
|
||||
return 0, fmt.Errorf("params.%s must be finite", key)
|
||||
}
|
||||
return v, nil
|
||||
}
|
||||
|
||||
func numberFromAny(raw any) (float64, bool) {
|
||||
switch v := raw.(type) {
|
||||
case float64:
|
||||
return v, true
|
||||
case float32:
|
||||
return float64(v), true
|
||||
case int:
|
||||
return float64(v), true
|
||||
case int8:
|
||||
return float64(v), true
|
||||
case int16:
|
||||
return float64(v), true
|
||||
case int32:
|
||||
return float64(v), true
|
||||
case int64:
|
||||
return float64(v), true
|
||||
case uint:
|
||||
return float64(v), true
|
||||
case uint8:
|
||||
return float64(v), true
|
||||
case uint16:
|
||||
return float64(v), true
|
||||
case uint32:
|
||||
return float64(v), true
|
||||
case uint64:
|
||||
return float64(v), true
|
||||
case string:
|
||||
parsed, err := strconv.ParseFloat(strings.TrimSpace(v), 64)
|
||||
return parsed, err == nil
|
||||
default:
|
||||
return 0, false
|
||||
}
|
||||
}
|
||||
|
||||
func latestIssueTime(raw []byte) time.Time {
|
||||
collection, err := spcprovider.DecodeGeoJSON(raw)
|
||||
if err != nil {
|
||||
return time.Time{}
|
||||
}
|
||||
var latest time.Time
|
||||
for _, feature := range collection.Features {
|
||||
t, err := spcprovider.ParseISOTimestamp(feature.Properties.IssueISO)
|
||||
if err != nil {
|
||||
continue
|
||||
}
|
||||
t = t.UTC()
|
||||
if latest.IsZero() || t.After(latest) {
|
||||
latest = t
|
||||
}
|
||||
}
|
||||
return latest
|
||||
}
|
||||
|
||||
func discussionUpdatedTime(rawHTML string) time.Time {
|
||||
text, err := spcprovider.ExtractProductText(rawHTML)
|
||||
if err != nil {
|
||||
return time.Time{}
|
||||
}
|
||||
t := spcprovider.ParseUpdatedTimestamp(text)
|
||||
if t == nil {
|
||||
return time.Time{}
|
||||
}
|
||||
return t.UTC()
|
||||
}
|
||||
|
||||
func chooseEffectiveTime(issue time.Time, updated time.Time, rss time.Time, fetched time.Time) time.Time {
|
||||
switch {
|
||||
case !issue.IsZero():
|
||||
return issue.UTC()
|
||||
case !updated.IsZero():
|
||||
return updated.UTC()
|
||||
case !rss.IsZero():
|
||||
return rss.UTC()
|
||||
default:
|
||||
return fetched.UTC()
|
||||
}
|
||||
}
|
||||
|
||||
func addHashPart(hash interface{ Write([]byte) (int, error) }, key, url string, body []byte) {
|
||||
_, _ = hash.Write([]byte(key))
|
||||
_, _ = hash.Write([]byte{0})
|
||||
_, _ = hash.Write([]byte(url))
|
||||
_, _ = hash.Write([]byte{0})
|
||||
_, _ = hash.Write(body)
|
||||
_, _ = hash.Write([]byte{0})
|
||||
}
|
||||
Reference in New Issue
Block a user