Implemented NWS weather stories support
All checks were successful
ci/woodpecker/push/build-image Pipeline was successful
All checks were successful
ci/woodpecker/push/build-image Pipeline was successful
This commit is contained in:
@@ -8,11 +8,15 @@
|
||||
// Canonical input schemas:
|
||||
// - weather.observation.v1 -> model.WeatherObservation
|
||||
// - weather.forecast.v1 -> model.WeatherForecastRun
|
||||
// - weather.forecast_discussion.v1 -> model.WeatherForecastDiscussion
|
||||
// - weather.weather_story.v1 -> model.WeatherStoryRun
|
||||
// - weather.alert.v1 -> model.WeatherAlertRun
|
||||
//
|
||||
// Parent/child relationships:
|
||||
// - observations.event_id -> observation_present_weather.event_id
|
||||
// - forecasts.event_id -> forecast_periods.run_event_id
|
||||
// - forecast_discussions.event_id -> forecast_discussion_key_messages.run_event_id
|
||||
// - weather_story_runs.event_id -> weather_stories.run_event_id
|
||||
// - alert_runs.event_id -> alerts.run_event_id
|
||||
// - alerts.(run_event_id, alert_index) -> alert_references.(run_event_id, alert_index)
|
||||
//
|
||||
@@ -24,13 +28,18 @@
|
||||
// - observation_present_weather.observed_at
|
||||
// - forecasts.issued_at
|
||||
// - forecast_periods.issued_at
|
||||
// - forecast_discussions.issued_at
|
||||
// - forecast_discussion_key_messages.issued_at
|
||||
// - weather_story_runs.as_of
|
||||
// - weather_stories.as_of
|
||||
// - alert_runs.as_of
|
||||
// - alerts.as_of
|
||||
// - alert_references.as_of
|
||||
//
|
||||
// Envelope field mapping (shared parent columns)
|
||||
//
|
||||
// These columns exist on observations, forecasts, and alert_runs:
|
||||
// These columns exist on parent tables such as observations, forecasts,
|
||||
// forecast_discussions, weather_story_runs, and alert_runs:
|
||||
// - event_id TEXT -> event.id
|
||||
// - event_kind TEXT -> event.kind
|
||||
// - event_source TEXT -> event.source
|
||||
@@ -120,7 +129,35 @@
|
||||
// - snowfall_depth_mm DOUBLE PRECISION NULL -> payload.periods[i].snowfallDepthMm
|
||||
// - uv_index DOUBLE PRECISION NULL -> payload.periods[i].uvIndex
|
||||
//
|
||||
// 5. alert_runs (PK: event_id)
|
||||
// 5. weather_story_runs (PK: event_id)
|
||||
//
|
||||
// - event_id TEXT -> event.id
|
||||
// - event_kind TEXT -> event.kind
|
||||
// - event_source TEXT -> event.source
|
||||
// - event_schema TEXT -> event.schema
|
||||
// - event_emitted_at TIMESTAMPTZ -> event.emitted_at
|
||||
// - event_effective_at TIMESTAMPTZ NULL -> event.effective_at
|
||||
// - office_id TEXT NULL -> payload.officeId
|
||||
// - as_of TIMESTAMPTZ -> payload.asOf
|
||||
// - story_count INTEGER -> len(payload.stories)
|
||||
//
|
||||
// 6. weather_stories (PK: run_event_id, story_index)
|
||||
//
|
||||
// - run_event_id TEXT -> weather_story_runs.event_id / payload.stories[i]
|
||||
// - story_index INTEGER -> i (array position in payload.stories)
|
||||
// - as_of TIMESTAMPTZ -> payload.asOf (copied from parent)
|
||||
// - office_id TEXT NULL -> payload.stories[i].officeId
|
||||
// - start_time TIMESTAMPTZ -> payload.stories[i].startTime
|
||||
// - end_time TIMESTAMPTZ -> payload.stories[i].endTime
|
||||
// - updated_at TIMESTAMPTZ -> payload.stories[i].updatedAt
|
||||
// - title TEXT NULL -> payload.stories[i].title
|
||||
// - description TEXT NULL -> payload.stories[i].description
|
||||
// - alt_text TEXT NULL -> payload.stories[i].altText
|
||||
// - priority BOOLEAN -> payload.stories[i].priority
|
||||
// - story_order INTEGER -> payload.stories[i].order
|
||||
// - download_url TEXT NULL -> payload.stories[i].downloadUrl
|
||||
//
|
||||
// 7. alert_runs (PK: event_id)
|
||||
//
|
||||
// - event_id TEXT -> event.id
|
||||
// - event_kind TEXT -> event.kind
|
||||
@@ -135,7 +172,7 @@
|
||||
// - longitude DOUBLE PRECISION NULL -> payload.longitude
|
||||
// - alert_count INTEGER -> len(payload.alerts)
|
||||
//
|
||||
// 6. alerts (PK: run_event_id, alert_index)
|
||||
// 8. alerts (PK: run_event_id, alert_index)
|
||||
//
|
||||
// - run_event_id TEXT -> alert_runs.event_id / payload.alerts[i]
|
||||
// - alert_index INTEGER -> i (array position in payload.alerts)
|
||||
@@ -160,7 +197,7 @@
|
||||
// - sender_name TEXT NULL -> payload.alerts[i].senderName
|
||||
// - reference_count INTEGER -> len(payload.alerts[i].references)
|
||||
//
|
||||
// 7. alert_references (PK: run_event_id, alert_index, reference_index)
|
||||
// 9. alert_references (PK: run_event_id, alert_index, reference_index)
|
||||
//
|
||||
// - run_event_id TEXT -> alert_runs.event_id / payload.alerts[i].references[j]
|
||||
// - alert_index INTEGER -> i (array position in payload.alerts)
|
||||
@@ -181,6 +218,10 @@
|
||||
// read one row from forecasts, then join forecast_periods by run_event_id
|
||||
// ordered by period_index to rebuild periods.
|
||||
//
|
||||
// - WeatherStoryRun:
|
||||
// read one row from weather_story_runs, then join weather_stories by
|
||||
// run_event_id ordered by story_index to rebuild stories.
|
||||
//
|
||||
// - WeatherAlertRun:
|
||||
// read one row from alert_runs, join alerts by run_event_id ordered by
|
||||
// alert_index, then join alert_references by (run_event_id, alert_index)
|
||||
|
||||
@@ -22,6 +22,8 @@ func mapPostgresEvent(_ context.Context, e fkevent.Event) ([]fksinks.PostgresWri
|
||||
return mapForecastEvent(e)
|
||||
case standards.SchemaWeatherForecastDiscussionV1:
|
||||
return mapForecastDiscussionEvent(e)
|
||||
case standards.SchemaWeatherStoryV1:
|
||||
return mapWeatherStoryEvent(e)
|
||||
case standards.SchemaWeatherAlertV1:
|
||||
return mapAlertEvent(e)
|
||||
default:
|
||||
@@ -218,6 +220,59 @@ func mapForecastDiscussionEvent(e fkevent.Event) ([]fksinks.PostgresWrite, error
|
||||
return writes, nil
|
||||
}
|
||||
|
||||
func mapWeatherStoryEvent(e fkevent.Event) ([]fksinks.PostgresWrite, error) {
|
||||
run, err := decodePayload[model.WeatherStoryRun](e.Payload)
|
||||
if err != nil {
|
||||
return nil, fmt.Errorf("decode weather story payload: %w", err)
|
||||
}
|
||||
if run.AsOf.IsZero() {
|
||||
return nil, fmt.Errorf("decode weather story payload: asOf is required")
|
||||
}
|
||||
|
||||
asOf := run.AsOf.UTC()
|
||||
writes := make([]fksinks.PostgresWrite, 0, 1+len(run.Stories))
|
||||
writes = append(writes, fksinks.PostgresWrite{
|
||||
Table: tableWeatherStoryRuns,
|
||||
Values: map[string]any{
|
||||
"event_id": e.ID,
|
||||
"event_kind": string(e.Kind),
|
||||
"event_source": e.Source,
|
||||
"event_schema": e.Schema,
|
||||
"event_emitted_at": e.EmittedAt.UTC(),
|
||||
"event_effective_at": nullableTime(e.EffectiveAt),
|
||||
"office_id": nullableString(run.OfficeID),
|
||||
"as_of": asOf,
|
||||
"story_count": len(run.Stories),
|
||||
},
|
||||
})
|
||||
|
||||
for i, story := range run.Stories {
|
||||
if story.StartTime.IsZero() || story.EndTime.IsZero() || story.UpdatedAt.IsZero() {
|
||||
return nil, fmt.Errorf("decode weather story payload: stories[%d] startTime/endTime/updatedAt are required", i)
|
||||
}
|
||||
writes = append(writes, fksinks.PostgresWrite{
|
||||
Table: tableWeatherStories,
|
||||
Values: map[string]any{
|
||||
"run_event_id": e.ID,
|
||||
"story_index": i,
|
||||
"as_of": asOf,
|
||||
"office_id": nullableString(story.OfficeID),
|
||||
"start_time": story.StartTime.UTC(),
|
||||
"end_time": story.EndTime.UTC(),
|
||||
"updated_at": story.UpdatedAt.UTC(),
|
||||
"title": nullableString(story.Title),
|
||||
"description": nullableString(story.Description),
|
||||
"alt_text": nullableString(story.AltText),
|
||||
"priority": story.Priority,
|
||||
"story_order": story.Order,
|
||||
"download_url": nullableString(story.DownloadURL),
|
||||
},
|
||||
})
|
||||
}
|
||||
|
||||
return writes, nil
|
||||
}
|
||||
|
||||
func mapAlertEvent(e fkevent.Event) ([]fksinks.PostgresWrite, error) {
|
||||
run, err := decodePayload[model.WeatherAlertRun](e.Payload)
|
||||
if err != nil {
|
||||
|
||||
@@ -193,6 +193,76 @@ func TestMapPostgresEventForecastDiscussionStructPayload(t *testing.T) {
|
||||
assertAllWritesIncludeAllColumns(t, writes)
|
||||
}
|
||||
|
||||
func TestMapPostgresEventWeatherStoryStructPayload(t *testing.T) {
|
||||
run := model.WeatherStoryRun{
|
||||
OfficeID: "LSX",
|
||||
AsOf: time.Date(2026, 5, 30, 9, 0, 34, 0, time.UTC),
|
||||
Stories: []model.WeatherStory{
|
||||
{
|
||||
OfficeID: "LSX",
|
||||
StartTime: time.Date(2026, 5, 30, 8, 46, 0, 0, time.UTC),
|
||||
EndTime: time.Date(2026, 5, 31, 11, 0, 0, 0, time.UTC),
|
||||
UpdatedAt: time.Date(2026, 5, 30, 9, 0, 34, 0, time.UTC),
|
||||
Title: "Several Chances for Rain Through Monday",
|
||||
Description: "Scattered showers and thunderstorms.",
|
||||
AltText: "This slide shows the forecast.",
|
||||
Priority: true,
|
||||
Order: 1,
|
||||
DownloadURL: "https://api.weather.gov/offices/LSX/weatherstories/download/story-1",
|
||||
},
|
||||
},
|
||||
}
|
||||
|
||||
writes, err := mapPostgresEvent(context.Background(), testEvent(standards.SchemaWeatherStoryV1, "weather_story", run))
|
||||
if err != nil {
|
||||
t.Fatalf("mapPostgresEvent() error = %v", err)
|
||||
}
|
||||
if len(writes) != 2 {
|
||||
t.Fatalf("mapPostgresEvent() writes len = %d, want 2", len(writes))
|
||||
}
|
||||
if writes[0].Table != tableWeatherStoryRuns {
|
||||
t.Fatalf("writes[0].Table = %q, want %q", writes[0].Table, tableWeatherStoryRuns)
|
||||
}
|
||||
if got := writes[0].Values["story_count"]; got != 1 {
|
||||
t.Fatalf("weather_story_runs story_count = %#v, want 1", got)
|
||||
}
|
||||
if writes[1].Table != tableWeatherStories {
|
||||
t.Fatalf("writes[1].Table = %q, want %q", writes[1].Table, tableWeatherStories)
|
||||
}
|
||||
if got := writes[1].Values["download_url"]; got != "https://api.weather.gov/offices/LSX/weatherstories/download/story-1" {
|
||||
t.Fatalf("weather_stories download_url = %#v", got)
|
||||
}
|
||||
if got := writes[1].Values["story_order"]; got != 1 {
|
||||
t.Fatalf("weather_stories story_order = %#v, want 1", got)
|
||||
}
|
||||
|
||||
assertAllWritesIncludeAllColumns(t, writes)
|
||||
}
|
||||
|
||||
func TestMapPostgresEventWeatherStoryRejectsMissingAsOf(t *testing.T) {
|
||||
_, err := mapPostgresEvent(context.Background(), testEvent(standards.SchemaWeatherStoryV1, "weather_story", model.WeatherStoryRun{}))
|
||||
if err == nil {
|
||||
t.Fatalf("mapPostgresEvent() error = nil, want missing asOf error")
|
||||
}
|
||||
if !strings.Contains(err.Error(), "asOf is required") {
|
||||
t.Fatalf("error = %q, want asOf context", err)
|
||||
}
|
||||
}
|
||||
|
||||
func TestMapPostgresEventWeatherStoryRejectsMissingStoryTimes(t *testing.T) {
|
||||
run := model.WeatherStoryRun{
|
||||
AsOf: time.Date(2026, 5, 30, 9, 0, 34, 0, time.UTC),
|
||||
Stories: []model.WeatherStory{{Title: "missing times"}},
|
||||
}
|
||||
_, err := mapPostgresEvent(context.Background(), testEvent(standards.SchemaWeatherStoryV1, "weather_story", run))
|
||||
if err == nil {
|
||||
t.Fatalf("mapPostgresEvent() error = nil, want missing story times error")
|
||||
}
|
||||
if !strings.Contains(err.Error(), "stories[0] startTime/endTime/updatedAt are required") {
|
||||
t.Fatalf("error = %q, want story time context", err)
|
||||
}
|
||||
}
|
||||
|
||||
func TestMapPostgresEventMapPayload(t *testing.T) {
|
||||
run := model.WeatherForecastRun{
|
||||
IssuedAt: time.Date(2026, 3, 16, 18, 0, 0, 0, time.UTC),
|
||||
|
||||
@@ -11,6 +11,8 @@ const (
|
||||
tableForecastPeriods = "forecast_periods"
|
||||
tableForecastDiscussions = "forecast_discussions"
|
||||
tableForecastDiscussionKeyMessages = "forecast_discussion_key_messages"
|
||||
tableWeatherStoryRuns = "weather_story_runs"
|
||||
tableWeatherStories = "weather_stories"
|
||||
tableAlertRuns = "alert_runs"
|
||||
tableAlerts = "alerts"
|
||||
tableAlertReferences = "alert_references"
|
||||
@@ -174,6 +176,51 @@ func PostgresSchema() fksinks.PostgresSchema {
|
||||
{Name: "idx_wf_discussion_message_issued_at", Columns: []string{"issued_at"}},
|
||||
},
|
||||
},
|
||||
{
|
||||
Name: tableWeatherStoryRuns,
|
||||
Columns: []fksinks.PostgresColumn{
|
||||
{Name: "event_id", Type: "TEXT", Nullable: false},
|
||||
{Name: "event_kind", Type: "TEXT", Nullable: false},
|
||||
{Name: "event_source", Type: "TEXT", Nullable: false},
|
||||
{Name: "event_schema", Type: "TEXT", Nullable: false},
|
||||
{Name: "event_emitted_at", Type: "TIMESTAMPTZ", Nullable: false},
|
||||
{Name: "event_effective_at", Type: "TIMESTAMPTZ", Nullable: true},
|
||||
{Name: "office_id", Type: "TEXT", Nullable: true},
|
||||
{Name: "as_of", Type: "TIMESTAMPTZ", Nullable: false},
|
||||
{Name: "story_count", Type: "INTEGER", Nullable: false},
|
||||
},
|
||||
PrimaryKey: []string{"event_id"},
|
||||
PruneColumn: "as_of",
|
||||
Indexes: []fksinks.PostgresIndex{
|
||||
{Name: "idx_wf_story_run_office_as_of", Columns: []string{"office_id", "as_of"}},
|
||||
{Name: "idx_wf_story_run_as_of", Columns: []string{"as_of"}},
|
||||
},
|
||||
},
|
||||
{
|
||||
Name: tableWeatherStories,
|
||||
Columns: []fksinks.PostgresColumn{
|
||||
{Name: "run_event_id", Type: "TEXT REFERENCES weather_story_runs(event_id) ON DELETE CASCADE", Nullable: false},
|
||||
{Name: "story_index", Type: "INTEGER", Nullable: false},
|
||||
{Name: "as_of", Type: "TIMESTAMPTZ", Nullable: false},
|
||||
{Name: "office_id", Type: "TEXT", Nullable: true},
|
||||
{Name: "start_time", Type: "TIMESTAMPTZ", Nullable: false},
|
||||
{Name: "end_time", Type: "TIMESTAMPTZ", Nullable: false},
|
||||
{Name: "updated_at", Type: "TIMESTAMPTZ", Nullable: false},
|
||||
{Name: "title", Type: "TEXT", Nullable: true},
|
||||
{Name: "description", Type: "TEXT", Nullable: true},
|
||||
{Name: "alt_text", Type: "TEXT", Nullable: true},
|
||||
{Name: "priority", Type: "BOOLEAN", Nullable: false},
|
||||
{Name: "story_order", Type: "INTEGER", Nullable: false},
|
||||
{Name: "download_url", Type: "TEXT", Nullable: true},
|
||||
},
|
||||
PrimaryKey: []string{"run_event_id", "story_index"},
|
||||
PruneColumn: "as_of",
|
||||
Indexes: []fksinks.PostgresIndex{
|
||||
{Name: "idx_wf_stories_start_time", Columns: []string{"start_time"}},
|
||||
{Name: "idx_wf_stories_end_time", Columns: []string{"end_time"}},
|
||||
{Name: "idx_wf_stories_updated_at", Columns: []string{"updated_at"}},
|
||||
},
|
||||
},
|
||||
{
|
||||
Name: tableAlertRuns,
|
||||
Columns: []fksinks.PostgresColumn{
|
||||
|
||||
@@ -15,6 +15,8 @@ func TestWeatherPostgresSchemaShape(t *testing.T) {
|
||||
tableForecastPeriods: true,
|
||||
tableForecastDiscussions: true,
|
||||
tableForecastDiscussionKeyMessages: true,
|
||||
tableWeatherStoryRuns: true,
|
||||
tableWeatherStories: true,
|
||||
tableAlertRuns: true,
|
||||
tableAlerts: true,
|
||||
tableAlertReferences: true,
|
||||
@@ -40,3 +42,38 @@ func TestWeatherPostgresSchemaShape(t *testing.T) {
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
func TestWeatherPostgresSchemaIncludesWeatherStoryColumns(t *testing.T) {
|
||||
runColumns := columnsForTable(t, tableWeatherStoryRuns)
|
||||
if !runColumns["as_of"] {
|
||||
t.Fatalf("%s missing as_of column", tableWeatherStoryRuns)
|
||||
}
|
||||
if !runColumns["story_count"] {
|
||||
t.Fatalf("%s missing story_count column", tableWeatherStoryRuns)
|
||||
}
|
||||
|
||||
storyColumns := columnsForTable(t, tableWeatherStories)
|
||||
for _, col := range []string{"start_time", "end_time", "updated_at", "title", "description", "alt_text", "priority", "story_order", "download_url"} {
|
||||
if !storyColumns[col] {
|
||||
t.Fatalf("%s missing %s column", tableWeatherStories, col)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
func columnsForTable(t *testing.T, table string) map[string]bool {
|
||||
t.Helper()
|
||||
|
||||
schema := PostgresSchema()
|
||||
for _, tbl := range schema.Tables {
|
||||
if tbl.Name != table {
|
||||
continue
|
||||
}
|
||||
cols := make(map[string]bool, len(tbl.Columns))
|
||||
for _, col := range tbl.Columns {
|
||||
cols[col.Name] = true
|
||||
}
|
||||
return cols
|
||||
}
|
||||
t.Fatalf("missing table %q", table)
|
||||
return nil
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user