Update Postgres outlook storage for v2

This commit is contained in:
2026-06-12 04:29:00 +00:00
parent 21a35a5205
commit 4d2cddf801
4 changed files with 325 additions and 75 deletions

View File

@@ -27,7 +27,7 @@ func mapPostgresEvent(_ context.Context, e fkevent.Event) ([]fksinks.PostgresWri
return mapWeatherStoryEvent(e)
case standards.SchemaWeatherAlertV1:
return mapAlertEvent(e)
case standards.SchemaWeatherOutlookV1:
case standards.SchemaWeatherOutlookV2:
return mapOutlookEvent(e)
default:
return nil, nil
@@ -339,17 +339,22 @@ func mapOutlookEvent(e fkevent.Event) ([]fksinks.PostgresWrite, error) {
}
asOf := run.AsOf.UTC()
writes := make([]fksinks.PostgresWrite, 0, 1+len(run.Outlooks))
if err := validateOutlookDiscussions(run.Discussions); err != nil {
return nil, err
}
writes := make([]fksinks.PostgresWrite, 0, 1+len(run.Outlooks)+len(run.Discussions))
writes = append(writes, fksinks.PostgresWrite{
Table: tableOutlookRuns,
Values: parentEventValues(e, map[string]any{
"location_id": nullableString(run.LocationID),
"location_name": nullableString(run.LocationName),
"latitude": nullableFloat64(run.Latitude),
"longitude": nullableFloat64(run.Longitude),
"as_of": asOf,
"issued_at": nullableTime(run.IssuedAt),
"outlook_count": len(run.Outlooks),
"location_id": nullableString(run.LocationID),
"location_name": nullableString(run.LocationName),
"latitude": nullableFloat64(run.Latitude),
"longitude": nullableFloat64(run.Longitude),
"as_of": asOf,
"issued_at": nullableTime(run.IssuedAt),
"outlook_count": len(run.Outlooks),
"discussion_count": len(run.Discussions),
}),
})
@@ -381,9 +386,6 @@ func mapOutlookEvent(e fkevent.Event) ([]fksinks.PostgresWrite, error) {
"issued_at": outlook.IssuedAt.UTC(),
"expires_at": outlook.ExpiresAt.UTC(),
"forecaster": nullableString(outlook.Forecaster),
"headline": nullableString(""),
"summary": nullableString(""),
"discussion": nullableString(""),
"source_url": nullableString(outlook.SourceURL),
"image_url": nullableString(outlook.ImageURL),
"contains_location": outlook.ContainsLocation,
@@ -392,6 +394,22 @@ func mapOutlookEvent(e fkevent.Event) ([]fksinks.PostgresWrite, error) {
})
}
for i, discussion := range run.Discussions {
writes = append(writes, fksinks.PostgresWrite{
Table: tableOutlookDiscussions,
Values: map[string]any{
"run_event_id": e.ID,
"discussion_index": i,
"as_of": asOf,
"day": discussion.Day,
"headline": nullableString(discussion.Headline),
"summary": nullableString(discussion.Summary),
"discussion": nullableString(discussion.Discussion),
"updated_at": nullableTime(discussion.UpdatedAt),
},
})
}
return writes, nil
}
@@ -423,6 +441,28 @@ func validateOutlook(outlook model.WeatherOutlook, index int) error {
if len(outlook.Geometry) == 0 {
return fmt.Errorf("decode outlook payload: outlooks[%d].geometry is required", index)
}
if !outlook.ContainsLocation {
return fmt.Errorf("decode outlook payload: outlooks[%d].containsLocation must be true", index)
}
return nil
}
func validateOutlookDiscussions(discussions []model.WeatherOutlookDiscussion) error {
seenDays := map[int]int{}
for i, discussion := range discussions {
if discussion.Day < 1 || discussion.Day > 3 {
return fmt.Errorf("decode outlook payload: discussions[%d].day must be 1, 2, or 3", i)
}
if strings.TrimSpace(discussion.Headline) == "" &&
strings.TrimSpace(discussion.Summary) == "" &&
strings.TrimSpace(discussion.Discussion) == "" {
return fmt.Errorf("decode outlook payload: discussions[%d] headline, summary, or discussion is required", i)
}
if first, ok := seenDays[discussion.Day]; ok {
return fmt.Errorf("decode outlook payload: discussions[%d].day duplicates discussions[%d].day %d", i, first, discussion.Day)
}
seenDays[discussion.Day] = i
}
return nil
}