Compare commits
2 Commits
v0.5.3
...
cfe6748330
| Author | SHA1 | Date | |
|---|---|---|---|
| cfe6748330 | |||
| 5a1134b955 |
@@ -194,7 +194,8 @@ GET /conditions/current?format=json&precision=0
|
|||||||
GET /alerts/active
|
GET /alerts/active
|
||||||
```
|
```
|
||||||
|
|
||||||
Returns the latest stored alert run filtered to alerts active at request time.
|
Returns the latest stored alert run filtered to alerts active at request time,
|
||||||
|
omitting older alerts superseded by newer alert references in the same run.
|
||||||
|
|
||||||
Query parameters: `format`, `units`.
|
Query parameters: `format`, `units`.
|
||||||
|
|
||||||
@@ -216,6 +217,9 @@ at or before request time, and the alert end boundary is absent or after request
|
|||||||
time. The end boundary prefers `ends`; if `ends` is absent, `expires` is used as
|
time. The end boundary prefers `ends`; if `ends` is absent, `expires` is used as
|
||||||
a fallback for older rows or providers that do not supply an alert-period end.
|
a fallback for older rows or providers that do not supply an alert-period end.
|
||||||
`onset` is presented when available but is not used as the active boundary.
|
`onset` is presented when available but is not used as the active boundary.
|
||||||
|
After active-time filtering, alerts referenced by another alert in the same run
|
||||||
|
are omitted as superseded. References from update and cancel messages are both
|
||||||
|
honored, even when the referencing alert is not itself returned.
|
||||||
|
|
||||||
Alert fields include `id`, `event`, `headline`, `severity`, `urgency`,
|
Alert fields include `id`, `event`, `headline`, `severity`, `urgency`,
|
||||||
`certainty`, `status`, `messageType`, `category`, `response`, `description`,
|
`certainty`, `status`, `messageType`, `category`, `response`, `description`,
|
||||||
|
|||||||
@@ -136,7 +136,10 @@ behavior deterministic.
|
|||||||
the application service with `alertNow().UTC()` so active alert filtering uses
|
the application service with `alertNow().UTC()` so active alert filtering uses
|
||||||
the request-time instant while remaining deterministic in endpoint tests. The
|
the request-time instant while remaining deterministic in endpoint tests. The
|
||||||
application service prefers alert `ends` over `expires` when deciding whether an
|
application service prefers alert `ends` over `expires` when deciding whether an
|
||||||
alert has ended.
|
alert has ended. The application service also suppresses alerts referenced by
|
||||||
|
another alert in the same latest run. This supersession rule uses alert
|
||||||
|
references from update and cancel messages, even when the referencing message is
|
||||||
|
not returned by `/alerts/active`.
|
||||||
|
|
||||||
## Outlook Filters
|
## Outlook Filters
|
||||||
|
|
||||||
|
|||||||
56
docs/roadmap/current.md
Normal file
56
docs/roadmap/current.md
Normal file
@@ -0,0 +1,56 @@
|
|||||||
|
# Current Conditions Condition-Code Selection
|
||||||
|
|
||||||
|
## Summary
|
||||||
|
|
||||||
|
Improve `/conditions/current` so numeric conditions continue to aggregate from recent observations, but `conditionCode` is selected by source-balanced WMO family consensus instead of numeric maximum. The goal is to prevent a single bad high WMO code, such as an erroneous thunderstorm code, from dominating current conditions while still returning a useful code when providers report semantically similar conditions.
|
||||||
|
|
||||||
|
## Target Behavior
|
||||||
|
|
||||||
|
Current conditions continue to use the application observation window, currently `app.ObservationWindowMinutesDefault`.
|
||||||
|
|
||||||
|
Numeric and directional fields remain aggregate values over recent `observations` rows:
|
||||||
|
|
||||||
|
- temperature, apparent temperature, dewpoint, relative humidity, and wind speed use averages;
|
||||||
|
- wind direction uses circular averaging;
|
||||||
|
- `isDay` comes from the latest row in the window.
|
||||||
|
|
||||||
|
`conditionCode` uses source-balanced consensus:
|
||||||
|
|
||||||
|
1. Select the latest observation per `event_source` within the current window.
|
||||||
|
2. Each source contributes at most one WMO condition-code vote.
|
||||||
|
3. Map each voted WMO code to a semantic family.
|
||||||
|
4. Select the family with the highest source vote count.
|
||||||
|
5. If the family vote is tied, return `model.WMOUnknown`.
|
||||||
|
6. Within the winning family, select the most frequent exact WMO code.
|
||||||
|
7. If exact-code vote is tied within the winning family, select the first code by family-specific representative ranking.
|
||||||
|
8. If there are no recognized condition-code votes, return `model.WMOUnknown`.
|
||||||
|
|
||||||
|
Family mapping and tie ranking:
|
||||||
|
|
||||||
|
| Family | Codes / tie ranking |
|
||||||
|
| --- | --- |
|
||||||
|
| `clear_or_cloud` | `0`, `1`, `2`, `3` |
|
||||||
|
| `fog` | `45`, `48` |
|
||||||
|
| `drizzle` | `51`, `53`, `55`, `56`, `57` |
|
||||||
|
| `rain` | `61`, `63`, `65`, `80`, `81`, `82`, `66`, `67` |
|
||||||
|
| `snow` | `71`, `73`, `75`, `85`, `86`, `77` |
|
||||||
|
| `thunderstorm` | `95`, `96`, `99` |
|
||||||
|
|
||||||
|
Examples:
|
||||||
|
|
||||||
|
| Source votes | Result |
|
||||||
|
| --- | --- |
|
||||||
|
| `0`, `1`, `2` | `0` |
|
||||||
|
| `1`, `2`, `95` | `1` |
|
||||||
|
| `0`, `95` | `model.WMOUnknown` |
|
||||||
|
| `61`, `63`, `80` | `61` |
|
||||||
|
| `61`, `95`, `0` | `model.WMOUnknown` |
|
||||||
|
| only `95` | `95` |
|
||||||
|
|
||||||
|
## Policy Decisions
|
||||||
|
|
||||||
|
- No public response schema change is required.
|
||||||
|
- No weatherfeeder change or database migration is required because `observations.event_source`, `condition_code`, `observed_at`, and `event_emitted_at` already exist.
|
||||||
|
- Provider-specific blacklists, trust weights, and source priorities are intentionally out of scope for this change.
|
||||||
|
- `WMOUnknown` is preferable to falsely choosing between tied precipitation, thunderstorm, and clear/cloud families.
|
||||||
|
- Current conditions remain a latest-window read model, not a durable derived table.
|
||||||
110
docs/roadmap/implementation.md
Normal file
110
docs/roadmap/implementation.md
Normal file
@@ -0,0 +1,110 @@
|
|||||||
|
# Implement Source-Balanced Current Conditions Condition Codes
|
||||||
|
|
||||||
|
## Summary
|
||||||
|
|
||||||
|
Implement the target behavior in `docs/roadmap/current.md`: keep `/conditions/current` response shape unchanged, continue aggregating numeric observations over the current window, and replace SQL `MAX(condition_code)` with Go-based source-balanced WMO family consensus.
|
||||||
|
|
||||||
|
This is a behavior change only. Do not add public fields, change query parameters, alter weatherfeeder tables, or add provider-specific blacklists.
|
||||||
|
|
||||||
|
## Stage 1: Add Condition-Code Consensus Helpers
|
||||||
|
|
||||||
|
Add package-local helpers in `internal/adapters/outbound/postgres` for current-conditions condition-code selection.
|
||||||
|
|
||||||
|
Required data shape:
|
||||||
|
|
||||||
|
- define a small candidate row/type containing `event_source` and `condition_code`;
|
||||||
|
- the selector accepts latest-per-source candidates and returns `model.WMOCode`.
|
||||||
|
|
||||||
|
Required selector behavior:
|
||||||
|
|
||||||
|
- use the family mapping and tie ranking from `docs/roadmap/current.md` exactly;
|
||||||
|
- ignore unrecognized WMO codes for family voting;
|
||||||
|
- return `model.WMOUnknown` when no recognized candidates exist;
|
||||||
|
- count each source at most once;
|
||||||
|
- return `model.WMOUnknown` on tied family votes;
|
||||||
|
- inside the winning family, choose the most frequent exact code;
|
||||||
|
- when exact codes tie inside the winning family, choose by family ranking.
|
||||||
|
|
||||||
|
Add focused unit tests for the selector before wiring SQL changes.
|
||||||
|
|
||||||
|
## Stage 2: Split Current-Conditions Condition-Code Querying
|
||||||
|
|
||||||
|
Update the Postgres current-conditions read path so condition-code selection is no longer computed with `MAX(condition_code)`.
|
||||||
|
|
||||||
|
Required SQL behavior:
|
||||||
|
|
||||||
|
- keep the existing aggregate query for sample count, numeric averages, circular wind direction, and latest `is_day`;
|
||||||
|
- remove `MAX(condition_code)` from the aggregate query result;
|
||||||
|
- add a separate query that returns one latest condition-code candidate per `event_source` inside the same observation window;
|
||||||
|
- latest per source is ordered by `observed_at DESC, event_emitted_at DESC`;
|
||||||
|
- candidate columns should include `event_source` and `condition_code`; timestamp columns may remain SQL-only if used only for ordering.
|
||||||
|
|
||||||
|
Required repository flow:
|
||||||
|
|
||||||
|
1. query aggregate current conditions;
|
||||||
|
2. return `nil, nil` when the aggregate sample count maps to no data, preserving current behavior;
|
||||||
|
3. query condition-code candidates using the same observation window;
|
||||||
|
4. run the Go selector;
|
||||||
|
5. map numeric aggregate fields plus selected condition code into `app.CurrentConditions`.
|
||||||
|
|
||||||
|
Keep row DTOs and mappers local to the Postgres adapter. Do not move SQL or row types into `internal/app`.
|
||||||
|
|
||||||
|
## Stage 3: Preserve API And Presentation Behavior
|
||||||
|
|
||||||
|
Keep all existing `/conditions/current` API behavior except condition-code selection.
|
||||||
|
|
||||||
|
Required invariants:
|
||||||
|
|
||||||
|
- `conditionCode` remains present in JSON/XML/text responses through the existing presenter path;
|
||||||
|
- `conditionText` continues to derive from selected `conditionCode` and `isDay`;
|
||||||
|
- `format`, `units`, and `precision` behavior is unchanged;
|
||||||
|
- no `tz` support is added;
|
||||||
|
- no public response fields are added or removed.
|
||||||
|
|
||||||
|
Update existing endpoint or presenter tests only if expected condition codes need to change because of the new selector.
|
||||||
|
|
||||||
|
## Stage 4: Update Current-Behavior Documentation
|
||||||
|
|
||||||
|
After implementation is complete, update current-behavior docs outside roadmap:
|
||||||
|
|
||||||
|
- `docs/api.md`: current conditions aggregate numeric fields and choose `conditionCode` by source-balanced WMO family consensus;
|
||||||
|
- `docs/internal/postgres-repository.md`: current conditions use an aggregate row plus a latest-per-source condition-code candidate query;
|
||||||
|
- `docs/integrations/weatherfeeder-postgres.md`: current conditions use `observations.event_source` for source-balanced condition-code selection.
|
||||||
|
|
||||||
|
Do not document any unimplemented options such as provider blacklists, trust weights, or source priorities.
|
||||||
|
|
||||||
|
## Test Plan
|
||||||
|
|
||||||
|
Selector tests:
|
||||||
|
|
||||||
|
- `0`, `1`, `2` returns `0` by clear/cloud ranking;
|
||||||
|
- `1`, `2`, `95` returns `1` because clear/cloud family wins;
|
||||||
|
- `0`, `95` returns `model.WMOUnknown` because families tie;
|
||||||
|
- `61`, `63`, `80` returns `61` by rain ranking;
|
||||||
|
- `61`, `95`, `0` returns `model.WMOUnknown` because three families tie;
|
||||||
|
- only `95` returns `95`;
|
||||||
|
- unrecognized-only candidates return `model.WMOUnknown`;
|
||||||
|
- duplicate candidates from the same source do not produce multiple votes if the selector receives them.
|
||||||
|
|
||||||
|
Repository tests:
|
||||||
|
|
||||||
|
- current conditions no longer chooses the highest numeric condition code;
|
||||||
|
- latest condition-code candidate per source uses `observed_at DESC, event_emitted_at DESC`;
|
||||||
|
- aggregate no-sample behavior still returns nil;
|
||||||
|
- numeric aggregate fields, wind direction, and `isDay` mapping remain unchanged;
|
||||||
|
- SQL/query errors from the candidate query include operation context.
|
||||||
|
|
||||||
|
Verification commands:
|
||||||
|
|
||||||
|
```sh
|
||||||
|
go test ./internal/adapters/outbound/postgres ./internal/app
|
||||||
|
go test ./internal/adapters/inbound/httpapi ./internal/adapters/inbound/httpapi/presenter
|
||||||
|
go test ./...
|
||||||
|
```
|
||||||
|
|
||||||
|
## Assumptions
|
||||||
|
|
||||||
|
- `observations.event_source` is non-null in the weatherfeeder-owned schema and is safe to use as the source identity.
|
||||||
|
- The latest-per-source query can be implemented with existing Postgres features and does not require a migration.
|
||||||
|
- Current conditions remain repository-owned because they are a Postgres aggregate read model; no new app service method or public interface is needed.
|
||||||
|
- The implementation should remain standard-library and SQL based; do not add dependencies.
|
||||||
@@ -10,6 +10,8 @@ import (
|
|||||||
"gitea.maximumdirect.net/ejr/weatherfeeder/model"
|
"gitea.maximumdirect.net/ejr/weatherfeeder/model"
|
||||||
)
|
)
|
||||||
|
|
||||||
|
const nwsAlertURLPrefix = "https://api.weather.gov/alerts/"
|
||||||
|
|
||||||
// Repository defines outbound data access used by weatherapi use cases.
|
// Repository defines outbound data access used by weatherapi use cases.
|
||||||
type Repository interface {
|
type Repository interface {
|
||||||
LatestObservation(ctx context.Context) (*model.WeatherObservation, error)
|
LatestObservation(ctx context.Context) (*model.WeatherObservation, error)
|
||||||
@@ -77,9 +79,10 @@ func (s *Service) LatestActiveAlertRun(ctx context.Context, activeAt time.Time)
|
|||||||
}
|
}
|
||||||
|
|
||||||
out := cloneAlertRun(run)
|
out := cloneAlertRun(run)
|
||||||
|
supersededIDs := collectSupersededAlertIDs(out.Alerts)
|
||||||
alerts := out.Alerts[:0]
|
alerts := out.Alerts[:0]
|
||||||
for _, alert := range out.Alerts {
|
for _, alert := range out.Alerts {
|
||||||
if isActiveAlert(alert, activeAt) {
|
if isActiveAlert(alert, activeAt) && !isSupersededAlert(alert, supersededIDs) {
|
||||||
alerts = append(alerts, alert)
|
alerts = append(alerts, alert)
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
@@ -181,6 +184,41 @@ func isActiveAlert(alert model.WeatherAlert, activeAt time.Time) bool {
|
|||||||
return true
|
return true
|
||||||
}
|
}
|
||||||
|
|
||||||
|
func collectSupersededAlertIDs(alerts []model.WeatherAlert) map[string]struct{} {
|
||||||
|
supersededIDs := make(map[string]struct{})
|
||||||
|
for _, alert := range alerts {
|
||||||
|
for _, ref := range alert.References {
|
||||||
|
id := normalizeAlertID(referenceAlertID(ref))
|
||||||
|
if id != "" {
|
||||||
|
supersededIDs[id] = struct{}{}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
return supersededIDs
|
||||||
|
}
|
||||||
|
|
||||||
|
func isSupersededAlert(alert model.WeatherAlert, supersededIDs map[string]struct{}) bool {
|
||||||
|
id := normalizeAlertID(alert.ID)
|
||||||
|
if id == "" {
|
||||||
|
return false
|
||||||
|
}
|
||||||
|
_, ok := supersededIDs[id]
|
||||||
|
return ok
|
||||||
|
}
|
||||||
|
|
||||||
|
func referenceAlertID(ref model.AlertReference) string {
|
||||||
|
if strings.TrimSpace(ref.Identifier) != "" {
|
||||||
|
return ref.Identifier
|
||||||
|
}
|
||||||
|
return ref.ID
|
||||||
|
}
|
||||||
|
|
||||||
|
func normalizeAlertID(value string) string {
|
||||||
|
value = strings.TrimSpace(value)
|
||||||
|
value = strings.TrimPrefix(value, nwsAlertURLPrefix)
|
||||||
|
return value
|
||||||
|
}
|
||||||
|
|
||||||
func cloneOutlookRun(run *model.WeatherOutlookRun) *model.WeatherOutlookRun {
|
func cloneOutlookRun(run *model.WeatherOutlookRun) *model.WeatherOutlookRun {
|
||||||
out := *run
|
out := *run
|
||||||
out.Latitude = copyFloat64(run.Latitude)
|
out.Latitude = copyFloat64(run.Latitude)
|
||||||
|
|||||||
@@ -272,6 +272,60 @@ func TestServiceLatestActiveAlertRunUsesEndsBeforeExpires(t *testing.T) {
|
|||||||
assertAlertIDs(t, run, []string{"ends-after-active-expires-before", "expires-fallback"})
|
assertAlertIDs(t, run, []string{"ends-after-active-expires-before", "expires-fallback"})
|
||||||
}
|
}
|
||||||
|
|
||||||
|
func TestServiceLatestActiveAlertRunSuppressesReferencedOriginal(t *testing.T) {
|
||||||
|
activeAt := testTime(12)
|
||||||
|
original := testAlert("https://api.weather.gov/alerts/urn:oid:original", "Alert", testTimePtr(9), testTimePtr(10), nil, testTimePtr(13), testTimePtr(13))
|
||||||
|
update := testAlert("https://api.weather.gov/alerts/urn:oid:update", "Update", testTimePtr(11), testTimePtr(11), nil, testTimePtr(13), testTimePtr(13))
|
||||||
|
update.References = []model.AlertReference{{Identifier: "urn:oid:original"}}
|
||||||
|
unrelated := testAlert("https://api.weather.gov/alerts/urn:oid:unrelated", "Alert", testTimePtr(9), testTimePtr(10), nil, testTimePtr(13), testTimePtr(13))
|
||||||
|
repo := &fakeRepository{alerts: testAlertRunWithAlerts([]model.WeatherAlert{original, update, unrelated})}
|
||||||
|
svc := NewService(repo)
|
||||||
|
|
||||||
|
run, err := svc.LatestActiveAlertRun(context.Background(), activeAt)
|
||||||
|
if err != nil {
|
||||||
|
t.Fatalf("unexpected error: %v", err)
|
||||||
|
}
|
||||||
|
assertAlertIDs(t, run, []string{
|
||||||
|
"https://api.weather.gov/alerts/urn:oid:update",
|
||||||
|
"https://api.weather.gov/alerts/urn:oid:unrelated",
|
||||||
|
})
|
||||||
|
}
|
||||||
|
|
||||||
|
func TestServiceLatestActiveAlertRunCancelSuppressesReferencedOriginal(t *testing.T) {
|
||||||
|
activeAt := testTime(12)
|
||||||
|
original := testAlert("https://api.weather.gov/alerts/urn:oid:original", "Alert", testTimePtr(9), testTimePtr(10), nil, testTimePtr(13), testTimePtr(13))
|
||||||
|
cancel := testAlert("https://api.weather.gov/alerts/urn:oid:cancel", "Cancel", testTimePtr(11), testTimePtr(11), nil, testTimePtr(13), testTimePtr(13))
|
||||||
|
cancel.References = []model.AlertReference{{Identifier: "urn:oid:original"}}
|
||||||
|
repo := &fakeRepository{alerts: testAlertRunWithAlerts([]model.WeatherAlert{original, cancel})}
|
||||||
|
svc := NewService(repo)
|
||||||
|
|
||||||
|
run, err := svc.LatestActiveAlertRun(context.Background(), activeAt)
|
||||||
|
if err != nil {
|
||||||
|
t.Fatalf("unexpected error: %v", err)
|
||||||
|
}
|
||||||
|
assertAlertIDs(t, run, []string{})
|
||||||
|
}
|
||||||
|
|
||||||
|
func TestServiceLatestActiveAlertRunReferenceIdentifierPrecedenceAndIDFallback(t *testing.T) {
|
||||||
|
activeAt := testTime(12)
|
||||||
|
fromID := testAlert("urn:oid:from-id", "Alert", testTimePtr(9), testTimePtr(10), nil, testTimePtr(13), testTimePtr(13))
|
||||||
|
fromIdentifier := testAlert("urn:oid:from-identifier", "Alert", testTimePtr(9), testTimePtr(10), nil, testTimePtr(13), testTimePtr(13))
|
||||||
|
idFallback := testAlert("urn:oid:id-fallback", "Alert", testTimePtr(9), testTimePtr(10), nil, testTimePtr(13), testTimePtr(13))
|
||||||
|
update := testAlert("urn:oid:update", "Update", testTimePtr(11), testTimePtr(11), nil, testTimePtr(13), testTimePtr(13))
|
||||||
|
update.References = []model.AlertReference{
|
||||||
|
{ID: "urn:oid:from-id", Identifier: "urn:oid:from-identifier"},
|
||||||
|
{ID: "urn:oid:id-fallback"},
|
||||||
|
}
|
||||||
|
repo := &fakeRepository{alerts: testAlertRunWithAlerts([]model.WeatherAlert{fromID, fromIdentifier, idFallback, update})}
|
||||||
|
svc := NewService(repo)
|
||||||
|
|
||||||
|
run, err := svc.LatestActiveAlertRun(context.Background(), activeAt)
|
||||||
|
if err != nil {
|
||||||
|
t.Fatalf("unexpected error: %v", err)
|
||||||
|
}
|
||||||
|
assertAlertIDs(t, run, []string{"urn:oid:from-id", "urn:oid:update"})
|
||||||
|
}
|
||||||
|
|
||||||
func TestServiceDelegatesLatestConvectiveOutlookRun(t *testing.T) {
|
func TestServiceDelegatesLatestConvectiveOutlookRun(t *testing.T) {
|
||||||
repo := &fakeRepository{outlookRun: testOutlookRun()}
|
repo := &fakeRepository{outlookRun: testOutlookRun()}
|
||||||
svc := NewService(repo)
|
svc := NewService(repo)
|
||||||
|
|||||||
Reference in New Issue
Block a user