diff --git a/.env.example b/.env.example index d80c0f3..d989771 100644 --- a/.env.example +++ b/.env.example @@ -15,7 +15,7 @@ MONGO_DATABASE=miti99bot # Comma-separated module list. Empty = load every module, including any module # added later — so list them explicitly when a deployment should only gain a new # module deliberately. -MODULES=util,misc,random,amlich,wordle,loldle,lol,stock,gold,coin,stats,monkeyd,sticker,alias,blacklist +MODULES=util,misc,random,amlich,wordle,loldle,lol,stock,gold,coin,stats,monkeyd,sticker,alias,blacklist,weather # Telegram user id for owner-only commands (renamed from BOT_OWNER_ID). OWNER_ID= # Comma-separated admin Telegram user ids (renamed from ADMIN_USER_IDS). diff --git a/README.md b/README.md index 2859e5e..5a1f849 100644 --- a/README.md +++ b/README.md @@ -22,7 +22,7 @@ Atlas via long polling and an in-process cron scheduler. | `alias` | `/alias ` save a replied message under a name, then send it back with `/insert `, bare `/`, or inline `@botname `; `/aliases` lists, `/unalias` deletes. See [docs/aliases.md](docs/aliases.md) | | `blacklist` | Per-topic text deny-list with whitelist exceptions: `/blacklist_add`, `/blacklist_del`, `/whitelist_add`, `/whitelist_del`, `/blacklist_rules` lists both, `/blacklist_check` judges a text, `/blacklist` does either, `/whitelist_rnd` picks a random exception. Passive — the bot never scans chat. See [docs/blacklist.md](docs/blacklist.md) | | `monkeyd` | `/monkeyd_crawl [font_size]` export a monkeydd.com novel as a PDF, `/monkeyd_tags ` list its tags as hashtags | -| `thoitiet` | Weather from Open-Meteo: `/thoitiet` for the next 6 hours hour by hour, `/thoitiethomnay` for current conditions and today, `/thoitietngaymai` for tomorrow, `/thoitiettuannay` for the next 7 days. Each takes `[location...]` (Vietnamese with or without diacritics, shorthand like `hcm`/`hn`, or a foreign city); the default is Ho Chi Minh City | +| `weather` | Weather from Open-Meteo: `/thoitiet` for the next 6 hours hour by hour, `/thoitiethomnay` for current conditions and today, `/thoitietngaymai` for tomorrow, `/thoitiettuannay` for the next 7 days. Each takes `[location...]` (Vietnamese with or without diacritics, shorthand like `hcm`/`hn`, or a foreign city); the default is Ho Chi Minh City. Flood alerts for Tân Thuận (Q.7): `/thuyvan` shows the 5-day tide-peak forecast at Phú An and Nhà Bè (Đài KTTV Nam Bộ bulletin), the rain forecast, and live VNDMS river gauges within 30 km; `/thuyvan_subscribe` / `/thuyvan_unsubscribe` toggle a 10:30 ICT push (retried at 12:30 if that run fails) sent only when a forecast peak reaches báo động I (1.40 m) or rain reaches 50 mm | Commands marked (admin) require a user ID in `ADMIN_IDS` or the owner; commands marked (owner) require `OWNER_ID`. Both kinds are hidden from `/help` and the diff --git a/cmd/server/main.go b/cmd/server/main.go index 20a2a1d..bbb6775 100644 --- a/cmd/server/main.go +++ b/cmd/server/main.go @@ -35,8 +35,8 @@ import ( "github.com/tiennm99/miti99bot/internal/modules/stats" "github.com/tiennm99/miti99bot/internal/modules/sticker" "github.com/tiennm99/miti99bot/internal/modules/stock" - "github.com/tiennm99/miti99bot/internal/modules/thoitiet" "github.com/tiennm99/miti99bot/internal/modules/util" + "github.com/tiennm99/miti99bot/internal/modules/weather" "github.com/tiennm99/miti99bot/internal/modules/wordle" "github.com/tiennm99/miti99bot/internal/server" "github.com/tiennm99/miti99bot/internal/storage" @@ -106,7 +106,7 @@ func factories() map[string]modules.Factory { sticker.CollectionName: sticker.New, "alias": alias.New, "blacklist": blacklist.New, - "thoitiet": thoitiet.New, + weather.CollectionName: weather.New, } } diff --git a/go.mod b/go.mod index 3609849..6b1d15d 100644 --- a/go.mod +++ b/go.mod @@ -4,6 +4,7 @@ go 1.26.5 require ( github.com/go-telegram/bot v1.20.0 + github.com/ledongthuc/pdf v0.0.0-20260907135840-6c8c28e0e8a0 github.com/robfig/cron/v3 v3.0.1 github.com/testcontainers/testcontainers-go v0.43.0 github.com/testcontainers/testcontainers-go/modules/mongodb v0.43.0 diff --git a/go.sum b/go.sum index 79147bd..dc6b15c 100644 --- a/go.sum +++ b/go.sum @@ -56,6 +56,8 @@ github.com/kr/pretty v0.3.1 h1:flRD4NNwYAUpkphVc1HcthR4KEIFJ65n8Mw5qdRn3LE= github.com/kr/pretty v0.3.1/go.mod h1:hoEshYVHaxMs3cyo3Yncou5ZscifuDolrwPKZanG3xk= github.com/kr/text v0.2.0 h1:5Nx0Ya0ZqY2ygV366QzturHI13Jq95ApcVaJBhpS+AY= github.com/kr/text v0.2.0/go.mod h1:eLer722TekiGuMkidMxC/pM04lWEeraHUUmBw8l2grE= +github.com/ledongthuc/pdf v0.0.0-20260907135840-6c8c28e0e8a0 h1:7Q+xNAZFmnfYOMweHN3c/PDFUKKfY1pVJ26K++QvVfU= +github.com/ledongthuc/pdf v0.0.0-20260907135840-6c8c28e0e8a0/go.mod h1:1fEHWurg7pvf5SG6XNE5Q8UZmOwex51Mkx3SLhrW5B4= github.com/lufia/plan9stats v0.0.0-20211012122336-39d0f177ccd0 h1:6E+4a0GO5zZEnZ81pIr0yLvtUWk2if982qA3F3QD6H4= github.com/lufia/plan9stats v0.0.0-20211012122336-39d0f177ccd0/go.mod h1:zJYVVT2jmtg6P3p1VtQj7WsuWi/y4VnjVBn7F8KPB3I= github.com/magiconair/properties v1.8.10 h1:s31yESBquKXCV9a/ScB3ESkOjUYYv+X0rg8SYxI99mE= diff --git a/internal/modules/thoitiet/api_client.go b/internal/modules/weather/api_client.go similarity index 86% rename from internal/modules/thoitiet/api_client.go rename to internal/modules/weather/api_client.go index 245b172..858e92d 100644 --- a/internal/modules/thoitiet/api_client.go +++ b/internal/modules/weather/api_client.go @@ -1,4 +1,4 @@ -package thoitiet +package weather import ( "context" @@ -164,21 +164,41 @@ func fetchForecast(ctx context.Context, client *http.Client, p place) (forecast, } func getJSON(ctx context.Context, client *http.Client, rawURL string, out any) error { - req, err := http.NewRequestWithContext(ctx, http.MethodGet, rawURL, nil) + body, err := getBody(ctx, client, rawURL, maxBodySize, http.Header{"Accept": {"application/json"}}) if err != nil { - return fmt.Errorf("build request: %w", err) + return err } - req.Header.Set("Accept", "application/json") - resp, err := client.Do(req) - if err != nil { - return fmt.Errorf("request: %w", err) - } - defer func() { _ = resp.Body.Close() }() - if resp.StatusCode != http.StatusOK { - return fmt.Errorf("status %d", resp.StatusCode) - } - if err := json.NewDecoder(io.LimitReader(resp.Body, maxBodySize)).Decode(out); err != nil { + if err := json.Unmarshal(body, out); err != nil { return fmt.Errorf("decode: %w", err) } return nil } + +// getBody GETs rawURL with the given headers and returns at most limit bytes +// of a 200 response. A longer body is an error rather than a silent cut, so a +// truncated JSON or PDF never reaches a parser. +func getBody(ctx context.Context, client *http.Client, rawURL string, limit int64, header http.Header) ([]byte, error) { + req, err := http.NewRequestWithContext(ctx, http.MethodGet, rawURL, nil) + if err != nil { + return nil, fmt.Errorf("build request: %w", err) + } + for k, v := range header { + req.Header[k] = v + } + resp, err := client.Do(req) + if err != nil { + return nil, fmt.Errorf("request: %w", err) + } + defer func() { _ = resp.Body.Close() }() + if resp.StatusCode != http.StatusOK { + return nil, fmt.Errorf("status %d", resp.StatusCode) + } + body, err := io.ReadAll(io.LimitReader(resp.Body, limit+1)) + if err != nil { + return nil, fmt.Errorf("read body: %w", err) + } + if int64(len(body)) > limit { + return nil, fmt.Errorf("body exceeds %d bytes", limit) + } + return body, nil +} diff --git a/internal/modules/weather/flood.go b/internal/modules/weather/flood.go new file mode 100644 index 0000000..f900206 --- /dev/null +++ b/internal/modules/weather/flood.go @@ -0,0 +1,270 @@ +package weather + +import ( + "context" + "fmt" + "net/http" + "strings" + "sync" + "time" +) + +// tanThuanPlace is the flood view's fixed location: phường Tân Thuận, Quận 7. +// The coordinates are Open-Meteo's geocoding result for "Tân Thuận" in Quận +// Bảy; geocoding the name at runtime returns 14 Vietnamese places. +var tanThuanPlace = place{ + Name: "Tân Thuận", + Latitude: 10.74111, + Longitude: 106.71806, + CountryCode: "VN", + Country: "Việt Nam", + Admin1: "Thành phố Hồ Chí Minh", +} + +const ( + // floodDays is the window the flood view and alert cover, starting today; + // it matches the tide bulletin's 5-day forecast. + floodDays = 5 + // nearbyRadiusKm bounds the live gauge list around Tân Thuận. Phú An + // (2.9 km) and Nhà Bè (6.8 km) are the gauges that matter; 30 km adds + // the upstream Sài Gòn and Đồng Nai stations. + nearbyRadiusKm = 30 + // heavyRainMm is the daily rain forecast that triggers an alert even at + // low tide. + heavyRainMm = 50 + // floodFetchTimeout covers the bulletin chain (homepage, article, PDF). + floodFetchTimeout = 15 * time.Second + // maxBulletinAge is how old the newest linked bulletin may be. Before the + // day's bulletin is out the homepage still links yesterday's; anything + // older means the site stopped publishing, and its dates no longer cover + // the window. + maxBulletinAge = 24 * time.Hour + + floodSourceLine = "Nguồn: Đài KTTV Nam Bộ, VNDMS, Open-Meteo" +) + +// tideAlarmLevels are the báo động I/II/III thresholds, in metres, that KTTV +// sets for both Phú An and Nhà Bè. The bulletin prints them and VNDMS +// returns the same values. +var tideAlarmLevels = [alarmTiers]float64{1.40, 1.50, 1.60} + +// tideTier maps a water level to its báo động tier, 0 below BĐ I. +func tideTier(height float64) int { + tier := 0 + for i, level := range tideAlarmLevels { + if height >= level { + tier = i + 1 + } + } + return tier +} + +var tierNames = [...]string{"", "BĐ I", "BĐ II", "BĐ III"} + +// floodReport gathers the three sources. Each fails on its own, so the view +// can still show what it has. +type floodReport struct { + Today time.Time // start of today, ICT + Bulletin tideBulletin + BulletinErr error + Rain dailyWeather + RainErr error + Stations []nearbyStation + StationsErr error +} + +// fetchFloodReport queries the tide bulletin, the Tân Thuận rain forecast and +// the live gauges concurrently. +func fetchFloodReport(ctx context.Context, client *http.Client, now time.Time) floodReport { + r := floodReport{Today: startOfDay(now.In(ictZone))} + var wg sync.WaitGroup + wg.Go(func() { + r.Bulletin, r.BulletinErr = fetchTideBulletin(ctx, client) + if r.BulletinErr == nil && r.Today.Sub(r.Bulletin.Issued) > maxBulletinAge { + r.BulletinErr = fmt.Errorf("tide bulletin: newest is from %s", r.Bulletin.Issued.Format("2006-01-02")) + } + }) + wg.Go(func() { + var f forecast + f, r.RainErr = fetchForecast(ctx, client, tanThuanPlace) + r.Rain = f.Daily + }) + wg.Go(func() { + var all []gaugeStation + if all, r.StationsErr = fetchGaugeStations(ctx, client); r.StationsErr == nil { + r.Stations = stationsNear(all, tanThuanPlace, nearbyRadiusKm) + } + }) + wg.Wait() + return r +} + +func startOfDay(t time.Time) time.Time { + return time.Date(t.Year(), t.Month(), t.Day(), 0, 0, 0, 0, t.Location()) +} + +// floodDay is one day of the flood view. +type floodDay struct { + Date time.Time + Peaks []tidePeak // highest tide per bulletin station: Phú An, Nhà Bè + HasTide bool + Tier int // báo động tier of the higher peak + RainMm float64 + HasRain bool +} + +// atRisk reports whether the day should raise an alert. +func (d floodDay) atRisk() bool { return d.Tier > 0 || (d.HasRain && d.RainMm >= heavyRainMm) } + +// days builds the floodDays-day window from today, joining the bulletin's +// Phú An and Nhà Bè peaks with the rain forecast by date. A day either +// source lacks is still listed, with that part missing. +func (r floodReport) days() []floodDay { + out := make([]floodDay, floodDays) + for i := range out { + d := floodDay{Date: r.Today.AddDate(0, 0, i)} + if r.BulletinErr == nil { + d.Peaks, d.HasTide = r.Bulletin.peaksOn(d.Date) + for _, p := range d.Peaks { + d.Tier = max(d.Tier, tideTier(p.Height)) + } + } + if r.RainErr == nil { + d.RainMm, d.HasRain = r.Rain.precipitationOn(d.Date) + } + out[i] = d + } + return out +} + +// peaksOn returns the day's highest tide at Phú An and at Nhà Bè. +func (b tideBulletin) peaksOn(date time.Time) ([]tidePeak, bool) { + var peaks []tidePeak + for _, station := range b.Stations[:2] { + for _, d := range station { + if !d.Date.Equal(date) { + continue + } + if p, ok := d.maxPeak(); ok { + peaks = append(peaks, p) + } + } + } + return peaks, len(peaks) > 0 +} + +// precipitationOn returns the forecast rain total for date. +func (d dailyWeather) precipitationOn(date time.Time) (float64, bool) { + key := date.Format("2006-01-02") + for i := range min(len(d.Time), len(d.PrecipitationSum)) { + if d.Time[i] == key { + return d.PrecipitationSum[i], true + } + } + return 0, false +} + +// anySource reports whether at least one source answered. +func (r floodReport) anySource() bool { + return r.BulletinErr == nil || r.RainErr == nil || r.StationsErr == nil +} + +// formatFloodReport renders /thuyvan. +func formatFloodReport(r floodReport) string { + var b strings.Builder + b.WriteString("🌊 Thuỷ văn — Tân Thuận (Q.7)\n") + switch { + case r.BulletinErr == nil && r.RainErr != nil: + fmt.Fprintf(&b, "Chưa lấy được dự báo mưa.\nĐỉnh triều dự báo, Phú An / Nhà Bè (bản tin %s):\n", r.Bulletin.Issued.Format("02/01")) + case r.BulletinErr == nil: + fmt.Fprintf(&b, "Đỉnh triều dự báo, Phú An / Nhà Bè (bản tin %s):\n", r.Bulletin.Issued.Format("02/01")) + case r.RainErr == nil: + b.WriteString("Chưa lấy được bản tin dự báo triều.\nMưa dự báo:\n") + default: + b.WriteString("Chưa lấy được dự báo triều và mưa.\n") + } + if r.BulletinErr == nil || r.RainErr == nil { + for _, d := range r.days() { + b.WriteString(formatFloodDay(d) + "\n") + } + } + switch { + case r.StationsErr != nil: + b.WriteString("Chưa lấy được mực nước hiện tại.\n") + case len(r.Stations) > 0: + b.WriteString("Mực nước hiện tại:\n") + for _, s := range r.Stations { + b.WriteString(formatNearbyStation(s) + "\n") + } + } + fmt.Fprintf(&b, "Báo động I/II/III tại Phú An, Nhà Bè: %s / %s / %s m\n", + formatMetres(tideAlarmLevels[0]), formatMetres(tideAlarmLevels[1]), formatMetres(tideAlarmLevels[2])) + b.WriteString(floodSourceLine) + return b.String() +} + +// formatFloodAlert renders the push for the at-risk days, or "" when no day +// is at risk. +func formatFloodAlert(r floodReport) string { + var lines []string + for _, d := range r.days() { + if d.atRisk() { + lines = append(lines, formatFloodDay(d)) + } + } + if len(lines) == 0 { + return "" + } + return "⚠️ Cảnh báo ngập — Tân Thuận (Q.7)\n" + strings.Join(lines, "\n") + + "\nXem chi tiết: /thuyvan\n" + floodSourceLine +} + +// formatFloodDay renders one day, for example +// "⚠️ 29/09: triều 1.60 / 1.58 m (17:00), mưa 62 mm — BĐ III, mưa lớn". +func formatFloodDay(d floodDay) string { + var parts []string + if d.HasTide { + heights := make([]string, len(d.Peaks)) + top := d.Peaks[0] + for i, p := range d.Peaks { + heights[i] = formatMetres(p.Height) + if p.Height > top.Height { + top = p + } + } + parts = append(parts, fmt.Sprintf("triều %s m (%s)", strings.Join(heights, " / "), top.Time)) + } + if d.HasRain { + parts = append(parts, "mưa "+decimal(d.RainMm)+" mm") + } + if len(parts) == 0 { + parts = append(parts, "chưa có số liệu") + } + var warnings []string + if d.Tier > 0 { + warnings = append(warnings, tierNames[d.Tier]) + } + if d.HasRain && d.RainMm >= heavyRainMm { + warnings = append(warnings, "mưa lớn") + } + line := d.Date.Format("02/01") + ": " + strings.Join(parts, ", ") + if len(warnings) > 0 { + line = "⚠️ " + line + " — " + strings.Join(warnings, ", ") + } + return line +} + +// formatNearbyStation renders one live gauge, for example +// "⚠️ Thủ Dầu Một (Sài Gòn, 26.1 km): 1.56 m lúc 7h 02/10 — BĐ II". +func formatNearbyStation(s nearbyStation) string { + line := fmt.Sprintf("%s (%s, %s km): %s m lúc %dh %02d/%02d", + s.Name, s.River, decimal(s.DistanceKm), formatMetres(s.Level), s.Hour, s.Day, s.Month) + if s.Tier > 0 { + line = "⚠️ " + line + " — " + tierNames[s.Tier] + } + return line +} + +// formatMetres renders a water level to the centimetre with the Vietnamese +// comma separator, as KTTV publishes it. +func formatMetres(m float64) string { return strings.Replace(fmt.Sprintf("%.2f", m), ".", ",", 1) } diff --git a/internal/modules/weather/flood_push.go b/internal/modules/weather/flood_push.go new file mode 100644 index 0000000..fc0ba14 --- /dev/null +++ b/internal/modules/weather/flood_push.go @@ -0,0 +1,216 @@ +package weather + +import ( + "context" + "errors" + "fmt" + "net/http" + "sync" + "time" + + "github.com/go-telegram/bot" + "github.com/go-telegram/bot/models" + + "github.com/tiennm99/miti99bot/internal/log" + "github.com/tiennm99/miti99bot/internal/modules" + "github.com/tiennm99/miti99bot/internal/modules/util/chathelper" + "github.com/tiennm99/miti99bot/internal/modules/util/subscription" +) + +const ( + floodFetchErrorText = "Không lấy được dữ liệu thuỷ văn. Thử lại sau nhé." + + // floodPushCronName is the cron's registry and scheduler key; it must be + // unique across all modules' crons. + floodPushCronName = "thuyvan_flood_push" + // floodPushSchedule is UTC: 03:30 UTC == 10:30 ICT, after the tide + // bulletin is published (about 09:20 ICT). + floodPushSchedule = "30 3 * * *" + // floodRetryCronName and floodRetrySchedule rerun the push at 12:30 ICT. + // The day is claimed only when an alert goes out, so the retry sends + // nothing after a successful 10:30 push and covers a 10:30 run that + // failed or found the day's bulletin not yet published. + floodRetryCronName = "thuyvan_flood_push_retry" + floodRetrySchedule = "30 5 * * *" + // floodPushDateKey records the ICT date of the last claimed push, so a + // double fire (a rolling deploy briefly running two containers) sends + // once. + floodPushDateKey = "flood_push:last_date" +) + +// flood is the /thuyvan state: its HTTP client and the alert subscribers. +type flood struct { + client *http.Client + subscribers subscription.Store + pushDate subscription.DayStore + // subscribersMu serializes Get→mutate→Put on the single subscriber slot. + subscribersMu sync.Mutex + // nowFn lets tests pin the clock; nil means time.Now. + nowFn func() time.Time +} + +func (f *flood) now() time.Time { + if f.nowFn != nil { + return f.nowFn() + } + return time.Now() +} + +func (f *flood) commands() []modules.Command { + return []modules.Command{ + { + Name: "thuyvan", + Visibility: modules.VisibilityPublic, + Description: "Dự báo triều cường, mưa và mực nước gần Tân Thuận (Q.7)", + Handler: f.handleReport, + }, + { + Name: "thuyvan_subscribe", + Visibility: modules.VisibilityPublic, + Description: "Nhận cảnh báo ngập Tân Thuận lúc 10:30 khi có triều cường hoặc mưa lớn", + Handler: f.handleSubscribe, + }, + { + Name: "thuyvan_unsubscribe", + Visibility: modules.VisibilityPublic, + Description: "Ngừng nhận cảnh báo ngập Tân Thuận", + Handler: f.handleUnsubscribe, + }, + } +} + +func (f *flood) crons() []modules.Cron { + return []modules.Cron{ + {Name: floodPushCronName, Schedule: floodPushSchedule, Handler: f.pushHandler}, + {Name: floodRetryCronName, Schedule: floodRetrySchedule, Handler: f.pushHandler}, + } +} + +func (f *flood) handleReport(ctx context.Context, b *bot.Bot, update *models.Update) error { + msg := update.Message + if msg == nil { + return nil + } + fetchCtx, cancelFetch := chathelper.FetchContext(ctx) + fetchCtx, cancel := context.WithTimeout(fetchCtx, floodFetchTimeout) + r := fetchFloodReport(fetchCtx, f.client, f.now()) + cancel() + cancelFetch() + logFloodErrors("/thuyvan", r) + if !r.anySource() { + return chathelper.Reply(ctx, b, msg, floodFetchErrorText) + } + return chathelper.Reply(ctx, b, msg, formatFloodReport(r)) +} + +func (f *flood) handleSubscribe(ctx context.Context, b *bot.Bot, update *models.Update) error { + msg := update.Message + if msg == nil { + return nil + } + f.subscribersMu.Lock() + defer f.subscribersMu.Unlock() + added, err := subscription.Add(ctx, f.subscribers, msg.Chat.ID, msg.MessageThreadID) + if err != nil { + return err + } + if added { + return chathelper.Reply(ctx, b, msg, "✅ Đã đăng ký cảnh báo ngập Tân Thuận cho "+scopeVI(msg)+ + ". Bot nhắn lúc 10:30 khi 5 ngày tới có triều cường từ báo động I hoặc mưa từ 50 mm.") + } + return chathelper.Reply(ctx, b, msg, capitalizeScopeVI(msg)+" đã đăng ký rồi.") +} + +func (f *flood) handleUnsubscribe(ctx context.Context, b *bot.Bot, update *models.Update) error { + msg := update.Message + if msg == nil { + return nil + } + f.subscribersMu.Lock() + defer f.subscribersMu.Unlock() + removed, err := subscription.Remove(ctx, f.subscribers, msg.Chat.ID, msg.MessageThreadID) + if err != nil { + return err + } + if removed { + return chathelper.Reply(ctx, b, msg, "Đã huỷ cảnh báo ngập cho "+scopeVI(msg)+".") + } + return chathelper.Reply(ctx, b, msg, capitalizeScopeVI(msg)+" chưa đăng ký.") +} + +// scopeVI names where a subscription applies, in Vietnamese. +func scopeVI(msg *models.Message) string { + if msg.MessageThreadID != 0 { + return "topic này" + } + return "chat này" +} + +func capitalizeScopeVI(msg *models.Message) string { + if msg.MessageThreadID != 0 { + return "Topic này" + } + return "Chat này" +} + +func (f *flood) pushHandler(ctx context.Context, deps modules.Deps) error { + if deps.Bot == nil { + return errors.New("flood push: deps.Bot is nil (BuildOptions.Bot not wired)") + } + return runFloodPush(ctx, f, deps.Bot) +} + +// runFloodPush sends the alert to every subscriber when a day in the window +// is at risk, and sends nothing otherwise. The day is claimed only after the +// fetch, so a failed fetch leaves it unclaimed. +func runFloodPush(ctx context.Context, f *flood, sender subscription.Sender) error { + subs, err := subscription.List(ctx, f.subscribers) + if err != nil { + return fmt.Errorf("flood push: list subscribers: %w", err) + } + if len(subs) == 0 { + return nil + } + fetchCtx, cancel := context.WithTimeout(ctx, floodFetchTimeout) + r := fetchFloodReport(fetchCtx, f.client, f.now()) + cancel() + logFloodErrors("flood push", r) + if r.BulletinErr != nil && r.RainErr != nil { + return fmt.Errorf("flood push: no tide or rain forecast: %w", errors.Join(r.BulletinErr, r.RainErr)) + } + text := formatFloodAlert(r) + if text == "" { + // Tide is the main flood signal here; without the bulletin a quiet + // rain forecast proves nothing, so fail loudly and let the retry + // run try again. + if r.BulletinErr != nil { + return fmt.Errorf("flood push: no tide forecast, rain alone shows no risk: %w", r.BulletinErr) + } + log.Info("flood push: no risk in the next days, skipping") + return nil + } + day := r.Today.Format("2006-01-02") + won, err := subscription.ClaimDay(ctx, f.pushDate, floodPushDateKey, day) + if err != nil { + return fmt.Errorf("flood push: claim date: %w", err) + } + if !won { + log.Info("flood push: already pushed today, skipping", "date", day) + return nil + } + res, err := subscription.Fanout(ctx, "flood", f.subscribers, &f.subscribersMu, subs, sender, bot.SendMessageParams{Text: text}) + if err != nil { + return err + } + log.Info("flood push complete", "subscribers", len(subs), + "sent", res.Sent, "failed", res.Failed, "pruned", res.Pruned, "throttled", res.Throttled) + return nil +} + +func logFloodErrors(where string, r floodReport) { + for source, err := range map[string]error{"tide bulletin": r.BulletinErr, "rain": r.RainErr, "gauges": r.StationsErr} { + if err != nil { + log.Warn("flood source failed", "module", "weather", "where", where, "source", source, "err", err) + } + } +} diff --git a/internal/modules/weather/flood_test.go b/internal/modules/weather/flood_test.go new file mode 100644 index 0000000..4410137 --- /dev/null +++ b/internal/modules/weather/flood_test.go @@ -0,0 +1,352 @@ +package weather + +import ( + "context" + "errors" + "net/http" + "net/http/httptest" + "os" + "strings" + "sync" + "testing" + "time" + + "github.com/go-telegram/bot" + "github.com/go-telegram/bot/models" + + "github.com/tiennm99/miti99bot/internal/modules" + "github.com/tiennm99/miti99bot/internal/modules/util/subscription" + "github.com/tiennm99/miti99bot/internal/storage" + "github.com/tiennm99/miti99bot/internal/testutil" +) + +// floodNow is 2026-10-02 12:00 ICT, the day of the bulletin fixture. +func floodNow() time.Time { return time.Date(2026, 10, 2, 5, 0, 0, 0, time.UTC) } + +// fakeFloodUpstreams serves KTTV Nam Bộ, VNDMS and Open-Meteo from one +// httptest server. A non-zero status fails that upstream. +type fakeFloodUpstreams struct { + forecastBody string + bulletin string // issue date YYYYMMDD of the served PDF; default 20261002 + kttvnbStatus int + vndmsStatus int + forecastStatus int + vndmsReferers []string + mu sync.Mutex +} + +func stubFloodUpstreams(t *testing.T, f *fakeFloodUpstreams) { + t.Helper() + if f.bulletin == "" { + f.bulletin = "20261002" + } + pdfData, err := os.ReadFile("testdata/HCMC_TVHN_" + f.bulletin + ".pdf") + if err != nil { + t.Fatal(err) + } + lv0, _ := os.ReadFile("testdata/vndms_lv0.json") + lv2, _ := os.ReadFile("testdata/vndms_lv2.json") + server := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + fail := func(status int) bool { + if status != 0 { + http.Error(w, "unavailable", status) + return true + } + return false + } + switch { + case r.URL.Path == "/kttvnb/": + if !fail(f.kttvnbStatus) { + d := f.bulletin + _, _ = w.Write([]byte(``)) + } + case strings.HasPrefix(r.URL.Path, "/kttvnb/index.php/"): + _, _ = w.Write([]byte(``)) + case strings.HasSuffix(r.URL.Path, ".pdf"): + _, _ = w.Write(pdfData) + case r.URL.Path == "/vndms/water_level": + f.mu.Lock() + f.vndmsReferers = append(f.vndmsReferers, r.Header.Get("Referer")) + f.mu.Unlock() + if fail(f.vndmsStatus) { + return + } + switch r.URL.Query().Get("lv") { + case "0": + _, _ = w.Write(lv0) + case "2": + _, _ = w.Write(lv2) + default: + _, _ = w.Write([]byte(`{"type":"FeatureCollection","features":[]}`)) + } + case r.URL.Path == "/forecast": + if !fail(f.forecastStatus) { + _, _ = w.Write([]byte(f.forecastBody)) + } + default: + http.NotFound(w, r) + } + })) + t.Cleanup(server.Close) + origK, origV, origF := kttvnbURL, vndmsURL, forecastURL + kttvnbURL, vndmsURL, forecastURL = server.URL+"/kttvnb", server.URL+"/vndms", server.URL+"/forecast" + t.Cleanup(func() { kttvnbURL, vndmsURL, forecastURL = origK, origV, origF }) +} + +func newTestFlood() *flood { + col := storage.NewMemoryProvider().Collection(CollectionName) + return &flood{ + client: http.DefaultClient, + subscribers: storage.Typed[subscription.Doc](col), + pushDate: storage.Typed[subscription.DayDoc](col), + nowFn: floodNow, + } +} + +func installFlood(t *testing.T, f *flood) *testutil.RecordingBot { + t.Helper() + rb := testutil.NewRecordingBot(t) + reg := &modules.Registry{ + Modules: []modules.Module{{Name: CollectionName, Commands: f.commands()}}, + AllCommands: map[string]modules.Command{}, + } + for _, c := range f.commands() { + reg.AllCommands[c.Name] = c + } + modules.Install(rb.Bot, reg, modules.Auth{}) + return rb +} + +func TestThuyvan_RendersAllSources(t *testing.T) { + up := &fakeFloodUpstreams{forecastBody: forecastFixture} + stubFloodUpstreams(t, up) + rb := installFlood(t, newTestFlood()) + + want := "🌊 Thuỷ văn — Tân Thuận (Q.7)\n" + + "Đỉnh triều dự báo, Phú An / Nhà Bè (bản tin 02/10):\n" + + "02/10: triều 1,33 / 1,34 m (06:00), mưa 7,9 mm\n" + + "03/10: triều 1,18 / 1,19 m (07:00), mưa 7,0 mm\n" + + "04/10: triều 1,01 / 1,03 m (21:30), mưa 5,7 mm\n" + + "05/10: triều 0,84 / 0,88 m (22:30), mưa 13,9 mm\n" + + "06/10: triều 1,13 / 0,70 m (01:30), mưa 4,0 mm\n" + + "Mực nước hiện tại:\n" + + "Phú An (Sài Gòn, 2,9 km): 1,33 m lúc 7h 02/10\n" + + "Nhà Bè (Đồng Điền, 6,8 km): 1,18 m lúc 7h 02/10\n" + + "Biên Hòa (Đồng Nai, 25,0 km): 1,74 m lúc 7h 02/10\n" + + "⚠️ Thủ Dầu Một (Sài Gòn, 26,1 km): 1,56 m lúc 7h 02/10 — BĐ II\n" + + "Báo động I/II/III tại Phú An, Nhà Bè: 1,40 / 1,50 / 1,60 m\n" + + floodSourceLine + if got := send(rb, "/thuyvan"); got != want { + t.Errorf("reply:\n%s\nwant:\n%s", got, want) + } + for _, ref := range up.vndmsReferers { + if ref != vndmsURL+"/" { + t.Errorf("VNDMS Referer = %q, want %q", ref, vndmsURL+"/") + } + } +} + +func TestThuyvan_PartialFailureKeepsOtherSources(t *testing.T) { + stubFloodUpstreams(t, &fakeFloodUpstreams{forecastBody: forecastFixture, kttvnbStatus: http.StatusBadGateway, vndmsStatus: http.StatusForbidden}) + rb := installFlood(t, newTestFlood()) + + got := send(rb, "/thuyvan") + for _, part := range []string{"Chưa lấy được bản tin dự báo triều.", "02/10: mưa 7,9 mm", "Chưa lấy được mực nước hiện tại."} { + if !strings.Contains(got, part) { + t.Errorf("reply missing %q:\n%s", part, got) + } + } + if strings.Contains(got, "triều 1,") { + t.Errorf("reply shows tide numbers without a bulletin:\n%s", got) + } +} + +func TestThuyvan_AllSourcesDown(t *testing.T) { + stubFloodUpstreams(t, &fakeFloodUpstreams{kttvnbStatus: 500, vndmsStatus: 500, forecastStatus: 500}) + rb := installFlood(t, newTestFlood()) + if got := send(rb, "/thuyvan"); got != floodFetchErrorText { + t.Errorf("reply = %q, want %q", got, floodFetchErrorText) + } +} + +func TestThuyvanSubscribe_Idempotent(t *testing.T) { + f := newTestFlood() + rb := installFlood(t, f) + if got := send(rb, "/thuyvan_subscribe"); !strings.Contains(got, "Đã đăng ký cảnh báo ngập") { + t.Errorf("first subscribe = %q", got) + } + if got := send(rb, "/thuyvan_subscribe"); got != "Chat này đã đăng ký rồi." { + t.Errorf("second subscribe = %q", got) + } + subs, _ := subscription.List(context.Background(), f.subscribers) + if len(subs) != 1 || subs[0].ChatID != 7 { + t.Errorf("subscribers = %+v, want chat 7", subs) + } + if got := send(rb, "/thuyvan_unsubscribe"); got != "Đã huỷ cảnh báo ngập cho chat này." { + t.Errorf("unsubscribe = %q", got) + } + if got := send(rb, "/thuyvan_unsubscribe"); got != "Chat này chưa đăng ký." { + t.Errorf("second unsubscribe = %q", got) + } +} + +type recordingSender struct { + calls []bot.SendMessageParams + err error +} + +func (s *recordingSender) SendMessage(_ context.Context, p *bot.SendMessageParams) (*models.Message, error) { + s.calls = append(s.calls, *p) + return &models.Message{}, s.err +} + +func TestRunFloodPush_SkipsWithoutRisk(t *testing.T) { + stubFloodUpstreams(t, &fakeFloodUpstreams{forecastBody: forecastFixture}) + f := newTestFlood() + _, _ = subscription.Add(context.Background(), f.subscribers, 7, 0) + sender := &recordingSender{} + if err := runFloodPush(context.Background(), f, sender); err != nil { + t.Fatal(err) + } + if len(sender.calls) != 0 { + t.Errorf("pushed %d messages on a day without risk", len(sender.calls)) + } + if _, _, err := f.pushDate.Get(context.Background(), floodPushDateKey); !errors.Is(err, storage.ErrNotFound) { + t.Errorf("a no-risk day must not claim the push date: err=%v", err) + } +} + +func TestRunFloodPush_HeavyRainAlertsOncePerDay(t *testing.T) { + // 62 mm forecast for 03/10. + stubFloodUpstreams(t, &fakeFloodUpstreams{forecastBody: strings.Replace(forecastFixture, + `"precipitation_sum":[5.30,7.90,7.00`, `"precipitation_sum":[5.30,7.90,62.00`, 1)}) + f := newTestFlood() + ctx := context.Background() + _, _ = subscription.Add(ctx, f.subscribers, 7, 0) + _, _ = subscription.Add(ctx, f.subscribers, 8, 3) + sender := &recordingSender{} + if err := runFloodPush(ctx, f, sender); err != nil { + t.Fatal(err) + } + if len(sender.calls) != 2 { + t.Fatalf("pushed %d messages, want 2", len(sender.calls)) + } + want := "⚠️ Cảnh báo ngập — Tân Thuận (Q.7)\n" + + "⚠️ 03/10: triều 1,18 / 1,19 m (07:00), mưa 62,0 mm — mưa lớn\n" + + "Xem chi tiết: /thuyvan\n" + floodSourceLine + if got := sender.calls[0].Text; got != want { + t.Errorf("push:\n%s\nwant:\n%s", got, want) + } + if sender.calls[1].ChatID != int64(8) || sender.calls[1].MessageThreadID != 3 { + t.Errorf("second push target = %v/%d, want 8/3", sender.calls[1].ChatID, sender.calls[1].MessageThreadID) + } + if err := runFloodPush(ctx, f, sender); err != nil { + t.Fatal(err) + } + if len(sender.calls) != 2 { + t.Errorf("second run the same day pushed again: %d messages", len(sender.calls)) + } +} + +func TestRunFloodPush_PrunesBlockedChats(t *testing.T) { + stubFloodUpstreams(t, &fakeFloodUpstreams{forecastBody: strings.Replace(forecastFixture, + `"precipitation_sum":[5.30,7.90`, `"precipitation_sum":[5.30,80.00`, 1)}) + f := newTestFlood() + ctx := context.Background() + _, _ = subscription.Add(ctx, f.subscribers, 7, 0) + sender := &recordingSender{err: errors.New("Forbidden: bot was blocked by the user")} + if err := runFloodPush(ctx, f, sender); err != nil { + t.Fatal(err) + } + if subs, _ := subscription.List(ctx, f.subscribers); len(subs) != 0 { + t.Errorf("blocked chat not pruned: %+v", subs) + } +} + +func TestRunFloodPush_FailsWhenNoForecast(t *testing.T) { + stubFloodUpstreams(t, &fakeFloodUpstreams{kttvnbStatus: 500, forecastStatus: 500}) + f := newTestFlood() + _, _ = subscription.Add(context.Background(), f.subscribers, 7, 0) + if err := runFloodPush(context.Background(), f, &recordingSender{}); err == nil { + t.Error("want an error when neither tide nor rain forecast is available") + } +} + +func TestThuyvan_RainFailureIsNoted(t *testing.T) { + stubFloodUpstreams(t, &fakeFloodUpstreams{forecastStatus: http.StatusBadGateway}) + rb := installFlood(t, newTestFlood()) + got := send(rb, "/thuyvan") + for _, part := range []string{"Chưa lấy được dự báo mưa.", "02/10: triều 1,33 / 1,34 m (06:00)\n"} { + if !strings.Contains(got, part) { + t.Errorf("reply missing %q:\n%s", part, got) + } + } +} + +func TestThuyvan_StaleBulletinIsNotUsed(t *testing.T) { + // The newest linked bulletin is from 10/09, three weeks before today. + stubFloodUpstreams(t, &fakeFloodUpstreams{forecastBody: forecastFixture, bulletin: "20260910"}) + rb := installFlood(t, newTestFlood()) + got := send(rb, "/thuyvan") + if !strings.Contains(got, "Chưa lấy được bản tin dự báo triều.") || strings.Contains(got, "bản tin 10/09") { + t.Errorf("stale bulletin used:\n%s", got) + } +} + +func TestRunFloodPush_TideAlert(t *testing.T) { + // The 10/09 bulletin forecasts báo động I–II at Phú An and Nhà Bè from + // 11/09; the rain fixture's dates do not overlap, so tide alone triggers. + stubFloodUpstreams(t, &fakeFloodUpstreams{forecastBody: forecastFixture, bulletin: "20260910"}) + f := newTestFlood() + f.nowFn = func() time.Time { return time.Date(2026, 9, 10, 4, 0, 0, 0, time.UTC) } + _, _ = subscription.Add(context.Background(), f.subscribers, 7, 0) + sender := &recordingSender{} + if err := runFloodPush(context.Background(), f, sender); err != nil { + t.Fatal(err) + } + if len(sender.calls) != 1 { + t.Fatalf("pushed %d messages, want 1", len(sender.calls)) + } + got := sender.calls[0].Text + for _, part := range []string{"⚠️ 11/09: triều 1,37 / 1,41 m (03:30) — BĐ I", "⚠️ 12/09: triều 1,45 / 1,46 m (16:00) — BĐ I"} { + if !strings.Contains(got, part) { + t.Errorf("push missing %q:\n%s", part, got) + } + } + if strings.Contains(got, "10/09") { + t.Errorf("a day below báo động I is in the push:\n%s", got) + } +} + +func TestRunFloodPush_FailsWithoutBulletinWhenRainIsQuiet(t *testing.T) { + stubFloodUpstreams(t, &fakeFloodUpstreams{forecastBody: forecastFixture, kttvnbStatus: 500}) + f := newTestFlood() + _, _ = subscription.Add(context.Background(), f.subscribers, 7, 0) + sender := &recordingSender{} + if err := runFloodPush(context.Background(), f, sender); err == nil { + t.Error("want an error so the missing tide forecast is not read as no risk") + } + if len(sender.calls) != 0 { + t.Errorf("pushed %d messages", len(sender.calls)) + } +} + +func TestFormatFloodDay_TideTiers(t *testing.T) { + cases := []struct { + height float64 + want string + }{ + {1.39, "29/09: triều 1,39 m (17:00)"}, + {1.40, "⚠️ 29/09: triều 1,40 m (17:00) — BĐ I"}, + {1.55, "⚠️ 29/09: triều 1,55 m (17:00) — BĐ II"}, + {1.60, "⚠️ 29/09: triều 1,60 m (17:00) — BĐ III"}, + } + for _, c := range cases { + d := floodDay{Date: time.Date(2026, 9, 29, 0, 0, 0, 0, ictZone), Peaks: []tidePeak{{c.height, "17:00"}}, + HasTide: true, Tier: tideTier(c.height)} + if got := formatFloodDay(d); got != c.want { + t.Errorf("height %.2f: %q, want %q", c.height, got, c.want) + } + } +} diff --git a/internal/modules/thoitiet/format.go b/internal/modules/weather/format.go similarity index 99% rename from internal/modules/thoitiet/format.go rename to internal/modules/weather/format.go index 22c0634..40cc240 100644 --- a/internal/modules/thoitiet/format.go +++ b/internal/modules/weather/format.go @@ -1,4 +1,4 @@ -package thoitiet +package weather import ( "fmt" diff --git a/internal/modules/thoitiet/location.go b/internal/modules/weather/location.go similarity index 99% rename from internal/modules/thoitiet/location.go rename to internal/modules/weather/location.go index e44e824..9275dfc 100644 --- a/internal/modules/thoitiet/location.go +++ b/internal/modules/weather/location.go @@ -1,4 +1,4 @@ -package thoitiet +package weather import ( "strings" diff --git a/internal/modules/thoitiet/location_test.go b/internal/modules/weather/location_test.go similarity index 99% rename from internal/modules/thoitiet/location_test.go rename to internal/modules/weather/location_test.go index 9d5871b..8b4dae1 100644 --- a/internal/modules/thoitiet/location_test.go +++ b/internal/modules/weather/location_test.go @@ -1,4 +1,4 @@ -package thoitiet +package weather import "testing" diff --git a/internal/modules/weather/testdata/HCMC_TVHN_20260910.pdf b/internal/modules/weather/testdata/HCMC_TVHN_20260910.pdf new file mode 100644 index 0000000..93e4bc3 Binary files /dev/null and b/internal/modules/weather/testdata/HCMC_TVHN_20260910.pdf differ diff --git a/internal/modules/weather/testdata/HCMC_TVHN_20261002.pdf b/internal/modules/weather/testdata/HCMC_TVHN_20261002.pdf new file mode 100644 index 0000000..39c538f Binary files /dev/null and b/internal/modules/weather/testdata/HCMC_TVHN_20261002.pdf differ diff --git a/internal/modules/weather/testdata/vndms_lv0.json b/internal/modules/weather/testdata/vndms_lv0.json new file mode 100644 index 0000000..b797737 --- /dev/null +++ b/internal/modules/weather/testdata/vndms_lv0.json @@ -0,0 +1 @@ +{"type":"FeatureCollection","features":[{"type":"Feature","geometry":{"type":"Point","coordinates":[106.63300323486328,17.433000564575195]},"properties":{"longtitude":106.633,"latitude":17.433,"label":"Đồng Hới","popupInfo":""}},{"type":"Feature","geometry":{"type":"Point","coordinates":[106.82499694824219,10.9399995803833]},"properties":{"longtitude":106.825,"latitude":10.94,"label":"Biên Hòa","popupInfo":"
  • Tên trạm: Biên Hòa
  • Mã trạm: 71594
  • Địa điểm: Đồng Nai
  • Sông: Đồng Nai
  • Nguồn: KTTV
  • Cảnh báo vượt mức báo động: -
  • Mực nước (1.74(m) 7-02/10)
  • Chi tiết
  • "}},{"type":"Feature","geometry":{"type":"Point","coordinates":[106.71700286865234,10.767000198364258]},"properties":{"longtitude":106.717,"latitude":10.767,"label":"Phú An","popupInfo":"
  • Tên trạm: Phú An
  • Mã trạm: 71600
  • Địa điểm: TP. Hồ Chí Minh
  • Sông: Sài Gòn
  • Nguồn: KTTV
  • Cảnh báo vượt mức báo động: -
  • Mực nước (1.33(m) 7-02/10)
  • Chi tiết
  • "}},{"type":"Feature","geometry":{"type":"Point","coordinates":[106.65299987792969,10.967000007629395]},"properties":{"longtitude":106.653,"latitude":10.967,"label":"Thủ Dầu Một","popupInfo":"
  • Tên trạm: Thủ Dầu Một
  • Mã trạm: 71586
  • Địa điểm: TP. Hồ Chí Minh
  • Sông: Sài Gòn
  • Nguồn: KTTV
  • Cảnh báo vượt mức báo động: -
  • Mực nước (1.56(m) 7-02/10)
  • Chi tiết
  • "}},{"type":"Feature","geometry":{"type":"Point","coordinates":[106.77100372314453,11.246000289916992]},"properties":{"longtitude":106.771,"latitude":11.246,"label":"Phước Hòa","popupInfo":"
  • Tên trạm: Phước Hòa
  • Mã trạm: 71585
  • Địa điểm: TP. Hồ Chí Minh
  • Sông: Bé
  • Nguồn: KTTV
  • Cảnh báo vượt mức báo động: -
  • Mực nước (26.12(m) 7-02/10)
  • Chi tiết
  • "}},{"type":"Feature","geometry":{"type":"Point","coordinates":[106.72000122070312,10.680000305175781]},"properties":{"longtitude":106.72,"latitude":10.68,"label":"Nhà Bè","popupInfo":"
  • Tên trạm: Nhà Bè
  • Mã trạm: 71601
  • Địa điểm: TP. Hồ Chí Minh
  • Sông: Đồng Điền
  • Nguồn: KTTV
  • Cảnh báo vượt mức báo động: -
  • Mực nước (1.18(m) 7-02/10)
  • Chi tiết
  • "}},{"type":"Feature","geometry":{"type":"Point","coordinates":[107.86666870117188,11.533333778381348]},"properties":{"longtitude":107.86667,"latitude":11.533334,"label":"Phước Hòa","popupInfo":"
  • Tên trạm: Phước Hòa
  • Địa điểm: TP. Hồ Chí Minh
  • Sông: Bé
  • Nguồn: KTTV
  • Cảnh báo vượt mức báo động: -
  • Mực nước :Không có số liệu
  • "}},{"type":"Feature","geometry":{"type":"Point","coordinates":[106.61666870117188,17.46666717529297]},"properties":{"longtitude":106.61667,"latitude":17.466667,"label":"Đồng Hới","popupInfo":"
  • Tên trạm: Đồng Hới
  • Địa điểm: Quảng Trị
  • Sông: Kiến Giang
  • Nguồn: KTTV
  • Cảnh báo vượt mức báo động: -
  • Mực nước :Không có số liệu
  • "}},{"type":"Feature","geometry":{"type":"Point","coordinates":[106.78333282470703,10.683333396911621]},"properties":{"longtitude":106.78333,"latitude":10.683333,"label":"Việt Lâm","popupInfo":"
  • Tên trạm: Việt Lâm
  • Địa điểm: Tuyên Quang
  • Sông: Lô
  • Nguồn: KTTV
  • Cảnh báo vượt mức báo động: -
  • Mực nước :Không có số liệu
  • "}}]} diff --git a/internal/modules/weather/testdata/vndms_lv2.json b/internal/modules/weather/testdata/vndms_lv2.json new file mode 100644 index 0000000..69a5d51 --- /dev/null +++ b/internal/modules/weather/testdata/vndms_lv2.json @@ -0,0 +1 @@ +{"type":"FeatureCollection","features":[{"type":"Feature","geometry":{"type":"Point","coordinates":[106.65299987792969,10.967000007629395]},"properties":{"longtitude":106.653,"latitude":10.967,"label":"Thủ Dầu Một","popupInfo":"
  • Tên trạm: Thủ Dầu Một
  • Mã trạm: 71586
  • Địa điểm: TP. Hồ Chí Minh
  • Sông: Sài Gòn
  • Nguồn: KTTV
  • Cảnh báo vượt mức báo động: -
  • Mực nước (1.56(m) 7-02/10)
  • Chi tiết
  • "}}]} diff --git a/internal/modules/weather/tide_bulletin.go b/internal/modules/weather/tide_bulletin.go new file mode 100644 index 0000000..4d2fd96 --- /dev/null +++ b/internal/modules/weather/tide_bulletin.go @@ -0,0 +1,365 @@ +package weather + +import ( + "bytes" + "context" + "errors" + "fmt" + "net/http" + "net/url" + "regexp" + "slices" + "strconv" + "strings" + "time" + + "github.com/ledongthuc/pdf" + "golang.org/x/text/unicode/norm" +) + +// ictZone is Vietnam's time zone; the bulletin's dates and times are local. +var ictZone = time.FixedZone("ICT", 7*60*60) + +// kttvnbURL is Đài KTTV Nam Bộ's site. Its homepage links each day's +// "bản tin dự báo thủy văn TPHCM" article, and the article links the 5-day +// tide bulletin PDF. A variable so tests can point it at an httptest server. +var kttvnbURL = "https://kttvnb.vn" + +const ( + // maxPageSize caps the homepage and article HTML (about 50 KB each). + maxPageSize = 1 << 20 + // maxBulletinSize caps the bulletin PDF (about 320 KB). + maxBulletinSize = 4 << 20 + // bulletinDays is how many forecast days each station block carries. + bulletinDays = 5 + // maxTideHeight bounds a plausible peak; anything outside (0, max] means + // the parser read the wrong column. + maxTideHeight = 3.0 +) + +// bulletinStations are the station blocks of the forecast table, in the +// order the PDF prints them. Only the first two are used, but all three are +// checked so a reordered table fails instead of mislabelling peaks. +var bulletinStations = []string{"Phú An", "Nhà Bè", "Thủ Dầu Một"} + +var ( + // articleLinkRE matches the homepage link to a TPHCM hydrology bulletin + // article, capturing the issue date. The slug varies + // ("...-tphcm-ra-nga-y-02-10-2026", "...-khu-va-c-tphcm-ra-nga-y-..."), so + // only its tail is matched. + articleLinkRE = regexp.MustCompile(`href="([^"]*tphcm-ra-nga-y-(\d{2})-(\d{2})-(\d{4}))"`) + // bulletinLinkRE matches the article's PDF attachment link. + bulletinLinkRE = regexp.MustCompile(`href="([^"]*/HCMC_TVHN_(\d{8})\.pdf)"`) + + dayTokenRE = regexp.MustCompile(`^\d{1,2}/$`) + monthTokenRE = regexp.MustCompile(`^\d{1,2}$`) + heightTokenRE = regexp.MustCompile(`^\d+\.\d+$`) + timeTokenRE = regexp.MustCompile(`^([01]?\d|2[0-3])\.[0-5]\d$`) +) + +// tidePeak is one forecast high tide. Time is "HH:MM" local time. +type tidePeak struct { + Height float64 + Time string +} + +// tideDay is one forecast day at one station: up to two high tides. +type tideDay struct { + Date time.Time + Peaks []tidePeak +} + +// maxPeak returns the day's highest tide, or false when the day has none. +func (d tideDay) maxPeak() (tidePeak, bool) { + if len(d.Peaks) == 0 { + return tidePeak{}, false + } + return slices.MaxFunc(d.Peaks, func(a, b tidePeak) int { + switch { + case a.Height < b.Height: + return -1 + case a.Height > b.Height: + return 1 + } + return 0 + }), true +} + +// tideBulletin is the parsed forecast table: one entry per station in +// bulletinStations order, each with bulletinDays days. +type tideBulletin struct { + Issued time.Time + Stations [][]tideDay +} + +// fetchTideBulletin finds the newest TPHCM bulletin linked from the KTTV Nam +// Bộ homepage and parses its forecast table. Before the day's bulletin is +// published (about 09:20 ICT) that is the previous day's. +func fetchTideBulletin(ctx context.Context, client *http.Client) (tideBulletin, error) { + home, err := getBody(ctx, client, kttvnbURL+"/", maxPageSize, nil) + if err != nil { + return tideBulletin{}, fmt.Errorf("tide bulletin homepage: %w", err) + } + articleURL, err := newestArticle(home) + if err != nil { + return tideBulletin{}, err + } + article, err := getBody(ctx, client, articleURL, maxPageSize, nil) + if err != nil { + return tideBulletin{}, fmt.Errorf("tide bulletin article: %w", err) + } + pdfURL, issued, err := bulletinLink(article) + if err != nil { + return tideBulletin{}, err + } + data, err := getBody(ctx, client, pdfURL, maxBulletinSize, nil) + if err != nil { + return tideBulletin{}, fmt.Errorf("tide bulletin pdf: %w", err) + } + b, err := parseTideBulletin(data, issued) + if err != nil { + return tideBulletin{}, fmt.Errorf("tide bulletin %s: %w", pdfURL, err) + } + return b, nil +} + +// newestArticle returns the absolute URL of the latest-dated bulletin article +// linked from the homepage. +func newestArticle(home []byte) (string, error) { + var best string + var bestDate time.Time + for _, m := range articleLinkRE.FindAllSubmatch(home, -1) { + d, err := time.Parse("02-01-2006", string(m[2])+"-"+string(m[3])+"-"+string(m[4])) + if err != nil || !d.After(bestDate) { + continue + } + best, bestDate = string(m[1]), d + } + if best == "" { + return "", errors.New("tide bulletin: no article link on homepage") + } + return resolveKTTVNB(best) +} + +// bulletinLink returns the article's PDF URL and the issue date encoded in +// its file name. +func bulletinLink(article []byte) (string, time.Time, error) { + m := bulletinLinkRE.FindSubmatch(article) + if m == nil { + return "", time.Time{}, errors.New("tide bulletin: no PDF link in article") + } + issued, err := time.ParseInLocation("20060102", string(m[2]), ictZone) + if err != nil { + return "", time.Time{}, fmt.Errorf("tide bulletin: issue date: %w", err) + } + u, err := resolveKTTVNB(string(m[1])) + return u, issued, err +} + +// resolveKTTVNB makes a possibly relative link absolute against kttvnbURL. +func resolveKTTVNB(link string) (string, error) { + base, err := url.Parse(kttvnbURL + "/") + if err != nil { + return "", err + } + ref, err := url.Parse(link) + if err != nil { + return "", fmt.Errorf("tide bulletin link %q: %w", link, err) + } + return base.ResolveReference(ref).String(), nil +} + +// parseTideBulletin extracts the forecast table from the bulletin PDF. +// +// The PDF's text layer garbles labels (letters split and reordered), but each +// forecast row keeps its date and values in reading order: +// +// 02/ 10 1.33 07.00 1.29 21.00 - 1.79 15.00 - 0.04 01.00 +// +// that is date, high 1 (height, time), high 2, low 1, low 2, with "ct" for a +// missing high. The date tokens can come out swapped ("10 04/"), with the +// station label in front on a block's middle row. Rows are matched by their +// date token, and the result is accepted only if it has bulletinDays +// consecutive dates for every station in bulletinStations, the station's +// label is an anagram of its name, and every height is plausible. The label +// sits in front of the block's middle row, or on its own text row a point +// away from it. +func parseTideBulletin(data []byte, issued time.Time) (b tideBulletin, err error) { + // The pdf package panics on some malformed files instead of returning an + // error; a bad download must not take the handler down with it. + defer func() { + if p := recover(); p != nil { + b, err = tideBulletin{}, fmt.Errorf("read pdf: %v", p) + } + }() + if !bytes.HasPrefix(data, []byte("%PDF-")) { + return tideBulletin{}, errors.New("not a PDF") + } + r, err := pdf.NewReader(bytes.NewReader(data), int64(len(data))) + if err != nil { + return tideBulletin{}, fmt.Errorf("open pdf: %w", err) + } + var rows []bulletinRow + var labels []textLine + for i := 1; i <= r.NumPage(); i++ { + page := r.Page(i) + if page.V.IsNull() { + continue + } + textRows, err := page.GetTextByRow() + if err != nil { + return tideBulletin{}, fmt.Errorf("read page %d: %w", i, err) + } + for _, tr := range textRows { + var tokens []string + for _, t := range tr.Content { + tokens = append(tokens, strings.Fields(t.S)...) + } + line := textLine{page: i, y: tr.Position, text: strings.Join(tokens, "")} + if row, ok := parseBulletinRow(tokens, issued); ok { + row.line = line + rows = append(rows, row) + } else { + labels = append(labels, line) + } + } + } + return assembleBulletin(rows, labels, issued) +} + +// textLine is one text row of the PDF and where it sits. +type textLine struct { + page int + y int64 + text string +} + +// labelRowGap is how far, in PDF points, a station label on its own row may +// sit from the forecast row it labels. +const labelRowGap = 3 + +// bulletinRow is one parsed forecast row plus the label text in front of it. +type bulletinRow struct { + label string + day tideDay + line textLine +} + +// parseBulletinRow reads one text row. It reports false for rows that are not +// forecast rows (headers, observed data, notes). +func parseBulletinRow(tokens []string, issued time.Time) (bulletinRow, bool) { + dayIdx := slices.IndexFunc(tokens, dayTokenRE.MatchString) + if dayIdx < 0 { + return bulletinRow{}, false + } + monthIdx := -1 + for _, i := range []int{dayIdx + 1, dayIdx - 1} { + if i >= 0 && i < len(tokens) && monthTokenRE.MatchString(tokens[i]) { + monthIdx = i + break + } + } + if monthIdx < 0 { + return bulletinRow{}, false + } + day, _ := strconv.Atoi(strings.TrimSuffix(tokens[dayIdx], "/")) + month, _ := strconv.Atoi(tokens[monthIdx]) + date, ok := bulletinDate(day, month, issued) + if !ok { + return bulletinRow{}, false + } + first := min(dayIdx, monthIdx) + values := tokens[max(dayIdx, monthIdx)+1:] + if len(values) < 4 { + return bulletinRow{}, false + } + var peaks []tidePeak + for i := 0; i < 4; i += 2 { + h, t := values[i], values[i+1] + if h == "ct" && t == "ct" { + continue + } + if !heightTokenRE.MatchString(h) || !timeTokenRE.MatchString(t) { + return bulletinRow{}, false + } + height, _ := strconv.ParseFloat(h, 64) + hour, minute, _ := strings.Cut(t, ".") + if len(hour) == 1 { + hour = "0" + hour + } + peaks = append(peaks, tidePeak{Height: height, Time: hour + ":" + minute}) + } + return bulletinRow{label: strings.Join(tokens[:first], ""), day: tideDay{Date: date, Peaks: peaks}}, true +} + +// bulletinDate builds a forecast date from the row's day and month. The table +// has no year: a month far before the issue month means the forecast crossed +// into the next year. +func bulletinDate(day, month int, issued time.Time) (time.Time, bool) { + if month < 1 || month > 12 || day < 1 || day > 31 { + return time.Time{}, false + } + year := issued.Year() + if month < int(issued.Month())-6 { + year++ + } + d := time.Date(year, time.Month(month), day, 0, 0, 0, 0, ictZone) + if d.Day() != day { + return time.Time{}, false + } + return d, true +} + +func assembleBulletin(rows []bulletinRow, labels []textLine, issued time.Time) (tideBulletin, error) { + want := len(bulletinStations) * bulletinDays + if len(rows) != want { + return tideBulletin{}, fmt.Errorf("found %d forecast rows, want %d", len(rows), want) + } + b := tideBulletin{Issued: issued} + for s, name := range bulletinStations { + block := rows[s*bulletinDays : (s+1)*bulletinDays] + if !labelled(block[bulletinDays/2], labels, name) { + return tideBulletin{}, fmt.Errorf("block %d is not labelled %s", s, name) + } + days := make([]tideDay, 0, bulletinDays) + for i, row := range block { + if !row.day.Date.Equal(rows[i].day.Date) || (i > 0 && !row.day.Date.Equal(days[i-1].Date.AddDate(0, 0, 1))) { + return tideBulletin{}, fmt.Errorf("%s: dates are not %d consecutive days", name, bulletinDays) + } + for _, p := range row.day.Peaks { + if p.Height <= 0 || p.Height > maxTideHeight { + return tideBulletin{}, fmt.Errorf("%s: implausible peak %.2f m", name, p.Height) + } + } + days = append(days, row.day) + } + b.Stations = append(b.Stations, days) + } + return b, nil +} + +// labelled reports whether row carries name as its label, either in front of +// its values or on a separate text row right next to it. +func labelled(row bulletinRow, labels []textLine, name string) bool { + if sameLetters(row.label, name) { + return true + } + for _, l := range labels { + gap := l.y - row.line.y + if l.page == row.line.page && gap >= -labelRowGap && gap <= labelRowGap && sameLetters(l.text, name) { + return true + } + } + return false +} + +// sameLetters reports whether a and b use the same letters, ignoring spaces +// and order; the PDF text layer scrambles a label's letters but keeps them. +func sameLetters(a, b string) bool { + letters := func(s string) []rune { + r := []rune(strings.Join(strings.Fields(norm.NFC.String(s)), "")) + slices.Sort(r) + return r + } + return slices.Equal(letters(a), letters(b)) +} diff --git a/internal/modules/weather/tide_bulletin_test.go b/internal/modules/weather/tide_bulletin_test.go new file mode 100644 index 0000000..4914b5e --- /dev/null +++ b/internal/modules/weather/tide_bulletin_test.go @@ -0,0 +1,130 @@ +package weather + +import ( + "os" + "testing" + "time" +) + +func loadBulletinFixture(t *testing.T) []byte { + t.Helper() + data, err := os.ReadFile("testdata/HCMC_TVHN_20261002.pdf") + if err != nil { + t.Fatal(err) + } + return data +} + +func TestParseTideBulletin_RealPDF(t *testing.T) { + issued := time.Date(2026, 10, 2, 0, 0, 0, 0, ictZone) + b, err := parseTideBulletin(loadBulletinFixture(t), issued) + if err != nil { + t.Fatalf("parseTideBulletin: %v", err) + } + if len(b.Stations) != len(bulletinStations) { + t.Fatalf("stations = %d, want %d", len(b.Stations), len(bulletinStations)) + } + phuAn, nhaBe := b.Stations[0], b.Stations[1] + if got := phuAn[0].Date; !got.Equal(issued) { + t.Errorf("first date = %v, want %v", got, issued) + } + want := []tidePeak{{1.33, "07:00"}, {1.29, "21:00"}} + if len(phuAn[0].Peaks) != 2 || phuAn[0].Peaks[0] != want[0] || phuAn[0].Peaks[1] != want[1] { + t.Errorf("Phú An 02/10 = %+v, want %+v", phuAn[0].Peaks, want) + } + // 05/10 has no second high tide ("ct"). + if got := phuAn[3].Peaks; len(got) != 1 || got[0] != (tidePeak{0.84, "10:00"}) { + t.Errorf("Phú An 05/10 = %+v, want one peak 0.84 m at 10:00", got) + } + // The middle row's date tokens come out swapped ("10 04/"). + if got := nhaBe[2]; !got.Date.Equal(issued.AddDate(0, 0, 2)) || got.Peaks[0] != (tidePeak{1.02, "08:00"}) { + t.Errorf("Nhà Bè 04/10 = %+v", got) + } + if p, _ := b.Stations[2][0].maxPeak(); p.Height != 1.60 { + t.Errorf("Thủ Dầu Một 02/10 max = %v, want 1.60", p.Height) + } +} + +// In the 10/09 bulletin the Thủ Dầu Một label is a text row of its own, one +// point from the block's middle row, and the week is at alarm level. +func TestParseTideBulletin_SeparateLabelRow(t *testing.T) { + data, err := os.ReadFile("testdata/HCMC_TVHN_20260910.pdf") + if err != nil { + t.Fatal(err) + } + b, err := parseTideBulletin(data, time.Date(2026, 9, 10, 0, 0, 0, 0, ictZone)) + if err != nil { + t.Fatalf("parseTideBulletin: %v", err) + } + if p, _ := b.Stations[1][2].maxPeak(); p != (tidePeak{1.46, "16:00"}) { + t.Errorf("Nhà Bè 12/09 max = %+v, want 1.46 m at 16:00", p) + } + if p, _ := b.Stations[2][2].maxPeak(); p != (tidePeak{1.55, "17:30"}) { + t.Errorf("Thủ Dầu Một 12/09 max = %+v, want 1.55 m at 17:30", p) + } +} + +func TestParseTideBulletin_RejectsNonBulletin(t *testing.T) { + if _, err := parseTideBulletin([]byte("not a pdf"), time.Now()); err == nil || err.Error() != "not a PDF" { + t.Error("want an error for non-PDF input") + } +} + +func TestParseBulletinRow(t *testing.T) { + issued := time.Date(2026, 12, 30, 0, 0, 0, 0, ictZone) + row, ok := parseBulletinRow([]string{"P", "n", "A", "ú", "h", "1", "02/", "1.01", "09.00", "ct", "ct", "1.78", "-"}, issued) + if !ok { + t.Fatal("row not parsed") + } + if want := time.Date(2027, 1, 2, 0, 0, 0, 0, ictZone); !row.day.Date.Equal(want) { + t.Errorf("date = %v, want %v (year rollover)", row.day.Date, want) + } + if !sameLetters(row.label, "Phú An") { + t.Errorf("label %q should match Phú An", row.label) + } + if len(row.day.Peaks) != 1 || row.day.Peaks[0] != (tidePeak{1.01, "09:00"}) { + t.Errorf("peaks = %+v, want one at 09:00", row.day.Peaks) + } + // A one-digit hour ("1.30") is padded. + if row, ok := parseBulletinRow([]string{"06/", "10", "0.68", "11.00", "1.13", "1.30"}, issued); !ok || row.day.Peaks[1].Time != "01:30" { + t.Errorf("one-digit hour: %+v ok=%v, want 01:30", row.day.Peaks, ok) + } + for _, tokens := range [][]string{ + {"đ", "ế", "n", "7h", "02/", "10"}, // header, no values + {"02/", "10", "61.56", "1538", "776", "256.0"}, // reservoir row, no times + {"MỰC", "NGÀY", "01/", "10/", "2026"}, // date without a month token + {"02/", "10", "1.33", "25.00", "1.29", "21.00"}, // hour out of range + } { + if _, ok := parseBulletinRow(tokens, issued); ok { + t.Errorf("parseBulletinRow(%q) should be rejected", tokens) + } + } +} + +func TestNewestArticleAndBulletinLink(t *testing.T) { + home := []byte(` +`) + got, err := newestArticle(home) + if err != nil { + t.Fatal(err) + } + if want := kttvnbURL + "/index.php/100-thong-tin-kttv/thuy-van/29194-ba-n-tin-da-ba-o-tha-y-v-n-tphcm-ra-nga-y-02-10-2026"; got != want { + t.Errorf("newestArticle = %q, want %q", got, want) + } + if _, err := newestArticle([]byte("")); err == nil { + t.Error("want an error when no article is linked") + } + + article := []byte(``) + pdfURL, issued, err := bulletinLink(article) + if err != nil { + t.Fatal(err) + } + if pdfURL != "https://kttvnb.vn/attachments/article/29194/HCMC_TVHN_20261002.pdf" { + t.Errorf("pdfURL = %q", pdfURL) + } + if !issued.Equal(time.Date(2026, 10, 2, 0, 0, 0, 0, ictZone)) { + t.Errorf("issued = %v", issued) + } +} diff --git a/internal/modules/weather/water_level.go b/internal/modules/weather/water_level.go new file mode 100644 index 0000000..4f4f7fc --- /dev/null +++ b/internal/modules/weather/water_level.go @@ -0,0 +1,178 @@ +package weather + +import ( + "context" + "encoding/json" + "fmt" + "html" + "math" + "net/http" + "regexp" + "sort" + "strconv" + "strings" + "sync" +) + +// vndmsURL is the Vietnam Disaster Monitoring System (Cục Quản lý đê điều và +// PCTT), which republishes KTTV river gauge readings. Its station-map +// endpoint is undocumented and answers 403 without a same-site Referer. A +// variable so tests can point it at an httptest server. +var vndmsURL = "https://vndms.gov.vn" + +const ( + // maxWaterLevelSize caps one station-map response; the full list is + // about 430 KB decoded. + maxWaterLevelSize = 4 << 20 + // alarmTiers is the number of báo động levels (I, II, III). The map + // endpoint's lv=N lists only the stations at tier N; lv=0 lists all. + alarmTiers = 3 +) + +var ( + popupNameRE = regexp.MustCompile(`Tên trạm: ([^<]+)`) + popupCodeRE = regexp.MustCompile(`Mã trạm: ([^<]+)`) + popupProvinceRE = regexp.MustCompile(`Địa điểm: ([^<]+)`) + popupRiverRE = regexp.MustCompile(`Sông: ([^<]+)`) + // popupLevelRE matches "Mực nước (1.33(m) 7-02/10)": metres, then the + // reading's hour and day/month. Stations without data print + // "Mực nước :Không có số liệu" instead and are skipped. + popupLevelRE = regexp.MustCompile(`Mực nước \((-?\d+(?:\.\d+)?)\(m\) (\d{1,2})-(\d{1,2})/(\d{1,2})\)`) +) + +// gaugeStation is one river gauge with its latest reading. +type gaugeStation struct { + Code, Name, River, Province string + Latitude, Longitude float64 + Level float64 // metres, in the station's own datum + Hour, Day, Month int // reading time, local + Tier int // báo động tier 1–3, or 0 below alarm +} + +type stationMap struct { + Features []struct { + Properties struct { + Latitude float64 `json:"latitude"` + Longitude float64 `json:"longtitude"` // sic + Popup string `json:"popupInfo"` + } `json:"properties"` + } `json:"features"` +} + +// fetchGaugeStations returns every station with a current reading, each +// tagged with its alarm tier. The four map requests run concurrently, and +// any failure fails the whole lookup so a missing tier list can never read as +// "below alarm". +func fetchGaugeStations(ctx context.Context, client *http.Client) ([]gaugeStation, error) { + maps := make([]stationMap, alarmTiers+1) + errs := make([]error, alarmTiers+1) + var wg sync.WaitGroup + for lv := range maps { + wg.Go(func() { + maps[lv], errs[lv] = fetchStationMap(ctx, client, lv) + }) + } + wg.Wait() + for lv, err := range errs { + if err != nil { + return nil, fmt.Errorf("water level lv=%d: %w", lv, err) + } + } + stations := parseStationMap(maps[0]) + if len(stations) == 0 { + return nil, fmt.Errorf("water level: no station readings") + } + tiers := map[string]int{} + for lv := 1; lv <= alarmTiers; lv++ { + for _, s := range parseStationMap(maps[lv]) { + tiers[s.Code] = max(tiers[s.Code], lv) + } + } + for i := range stations { + stations[i].Tier = tiers[stations[i].Code] + } + return stations, nil +} + +func fetchStationMap(ctx context.Context, client *http.Client, lv int) (stationMap, error) { + body, err := getBody(ctx, client, vndmsURL+"/water_level?lv="+strconv.Itoa(lv), maxWaterLevelSize, + http.Header{"Referer": {vndmsURL + "/"}, "Accept": {"application/json"}}) + if err != nil { + return stationMap{}, err + } + var m stationMap + if err := json.Unmarshal(body, &m); err != nil { + return stationMap{}, fmt.Errorf("decode: %w", err) + } + return m, nil +} + +// parseStationMap reads the station fields out of each feature's popup HTML, +// keeping only stations that have a code and a current reading. +func parseStationMap(m stationMap) []gaugeStation { + var out []gaugeStation + for _, f := range m.Features { + p := f.Properties.Popup + lm := popupLevelRE.FindStringSubmatch(p) + code := popupField(popupCodeRE, p) + if lm == nil || code == "" { + continue + } + level, err := strconv.ParseFloat(lm[1], 64) + if err != nil { + continue + } + hour, _ := strconv.Atoi(lm[2]) + day, _ := strconv.Atoi(lm[3]) + month, _ := strconv.Atoi(lm[4]) + out = append(out, gaugeStation{ + Code: code, + Name: popupField(popupNameRE, p), + River: popupField(popupRiverRE, p), + Province: popupField(popupProvinceRE, p), + Latitude: f.Properties.Latitude, + Longitude: f.Properties.Longitude, + Level: level, + Hour: hour, + Day: day, + Month: month, + }) + } + return out +} + +func popupField(re *regexp.Regexp, popup string) string { + m := re.FindStringSubmatch(popup) + if m == nil { + return "" + } + return html.UnescapeString(strings.TrimSpace(m[1])) +} + +// nearbyStation is a station with its distance from a reference point. +type nearbyStation struct { + gaugeStation + DistanceKm float64 +} + +// stationsNear returns the stations within radiusKm of p, nearest first. +func stationsNear(stations []gaugeStation, p place, radiusKm float64) []nearbyStation { + var out []nearbyStation + for _, s := range stations { + if d := distanceKm(p.Latitude, p.Longitude, s.Latitude, s.Longitude); d <= radiusKm { + out = append(out, nearbyStation{gaugeStation: s, DistanceKm: d}) + } + } + sort.SliceStable(out, func(i, j int) bool { return out[i].DistanceKm < out[j].DistanceKm }) + return out +} + +// distanceKm is the haversine great-circle distance. +func distanceKm(lat1, lon1, lat2, lon2 float64) float64 { + const earthRadiusKm = 6371 + rad := func(d float64) float64 { return d * math.Pi / 180 } + dLat, dLon := rad(lat2-lat1), rad(lon2-lon1) + a := math.Sin(dLat/2)*math.Sin(dLat/2) + + math.Cos(rad(lat1))*math.Cos(rad(lat2))*math.Sin(dLon/2)*math.Sin(dLon/2) + return 2 * earthRadiusKm * math.Asin(math.Sqrt(a)) +} diff --git a/internal/modules/weather/water_level_test.go b/internal/modules/weather/water_level_test.go new file mode 100644 index 0000000..6c47549 --- /dev/null +++ b/internal/modules/weather/water_level_test.go @@ -0,0 +1,59 @@ +package weather + +import ( + "encoding/json" + "os" + "testing" +) + +func loadStationMap(t *testing.T, name string) stationMap { + t.Helper() + data, err := os.ReadFile("testdata/" + name) + if err != nil { + t.Fatal(err) + } + var m stationMap + if err := json.Unmarshal(data, &m); err != nil { + t.Fatal(err) + } + return m +} + +func TestParseStationMap_SkipsStationsWithoutReadings(t *testing.T) { + stations := parseStationMap(loadStationMap(t, "vndms_lv0.json")) + // The fixture has 9 features; Việt Lâm and the duplicate Phước Hòa and + // Đồng Hới entries have no data, leaving 6 readings. + if len(stations) != 6 { + t.Fatalf("stations = %d, want 6: %+v", len(stations), stations) + } + var phuAn gaugeStation + for _, s := range stations { + if s.Name == "Việt Lâm" { + t.Error("Việt Lâm has no reading and should be skipped") + } + if s.Name == "Phú An" { + phuAn = s + } + } + want := gaugeStation{Code: "71600", Name: "Phú An", River: "Sài Gòn", Province: "TP. Hồ Chí Minh", + Latitude: phuAn.Latitude, Longitude: phuAn.Longitude, Level: 1.33, Hour: 7, Day: 2, Month: 10} + if phuAn != want { + t.Errorf("Phú An = %+v, want %+v", phuAn, want) + } +} + +func TestStationsNear_TanThuan(t *testing.T) { + near := stationsNear(parseStationMap(loadStationMap(t, "vndms_lv0.json")), tanThuanPlace, nearbyRadiusKm) + want := []string{"Phú An", "Nhà Bè", "Biên Hòa", "Thủ Dầu Một"} + if len(near) != len(want) { + t.Fatalf("near = %d stations, want %v", len(near), want) + } + for i, s := range near { + if s.Name != want[i] { + t.Errorf("near[%d] = %s, want %s", i, s.Name, want[i]) + } + } + if d := near[0].DistanceKm; d < 2.5 || d > 3.3 { + t.Errorf("Phú An distance = %.2f km, want about 2.9", d) + } +} diff --git a/internal/modules/thoitiet/thoitiet.go b/internal/modules/weather/weather.go similarity index 71% rename from internal/modules/thoitiet/thoitiet.go rename to internal/modules/weather/weather.go index 22928f8..f59aa8a 100644 --- a/internal/modules/thoitiet/thoitiet.go +++ b/internal/modules/weather/weather.go @@ -1,8 +1,12 @@ -// Package thoitiet is the weather module: /thoitiet shows the next 6 hours -// hour by hour, and /thoitiethomnay, /thoitietngaymai, and /thoitiettuannay -// show today's, tomorrow's, and the next 7 days' forecast for a location, Ho -// Chi Minh City by default. Data comes from Open-Meteo, which needs no API key. -package thoitiet +// Package weather is the weather and flood module. /thoitiet shows the next 6 +// hours hour by hour, and /thoitiethomnay, /thoitietngaymai, and +// /thoitiettuannay show today's, tomorrow's, and the next 7 days' forecast for +// a location, Ho Chi Minh City by default; that data comes from Open-Meteo, +// which needs no API key. /thuyvan shows the flood risk at Tân Thuận (Quận 7) +// from the KTTV Nam Bộ tide bulletin, the Open-Meteo rain forecast and VNDMS +// river gauges, and /thuyvan_subscribe opts a chat into a daily alert sent +// only when a tide or rain threshold is forecast. +package weather import ( "context" @@ -16,6 +20,8 @@ import ( "github.com/tiennm99/miti99bot/internal/log" "github.com/tiennm99/miti99bot/internal/modules" "github.com/tiennm99/miti99bot/internal/modules/util/chathelper" + "github.com/tiennm99/miti99bot/internal/modules/util/subscription" + "github.com/tiennm99/miti99bot/internal/storage" ) const ( @@ -46,9 +52,19 @@ var ( weekView = view{command: "thoitiettuannay", ready: hasDays(1), render: formatWeek} ) -// New is the thoitiet module Factory. The module keeps no state. -func New(_ modules.Deps) modules.Module { +// CollectionName is the module's registry key and storage collection, which +// holds the /thuyvan alert subscribers. +const CollectionName = "weather" + +// New is the weather module Factory. Only the flood alert keeps state: its +// subscriber list and last-push date. +func New(deps modules.Deps) modules.Module { client := &http.Client{Timeout: httpTimeout} + fl := &flood{ + client: &http.Client{Timeout: floodFetchTimeout}, + subscribers: storage.Typed[subscription.Doc](deps.Store), + pushDate: storage.Typed[subscription.DayDoc](deps.Store), + } command := func(name, description string, v view) modules.Command { return modules.Command{ Name: name, @@ -59,12 +75,13 @@ func New(_ modules.Deps) modules.Module { } } return modules.Module{ - Commands: []modules.Command{ + Commands: append([]modules.Command{ command("thoitiethomnay", "Thời tiết hôm nay (mặc định TP.HCM)", todayView), command("thoitiet", "Thời tiết từng giờ trong 6 giờ tới (mặc định TP.HCM)", hourlyView), command("thoitietngaymai", "Thời tiết ngày mai (mặc định TP.HCM)", tomorrowView), command("thoitiettuannay", "Thời tiết 7 ngày tới (mặc định TP.HCM)", weekView), - }, + }, fl.commands()...), + Crons: fl.crons(), } } @@ -85,7 +102,7 @@ func handler(client *http.Client, v view) modules.CommandHandler { err = fmt.Errorf("forecast: missing rows for /%s", v.command) } if err != nil { - log.Error("weather fetch failed", "module", "thoitiet", "command", v.command, "err", err) + log.Error("weather fetch failed", "module", "weather", "command", v.command, "err", err) return chathelper.Reply(ctx, b, msg, fetchErrorText) } return chathelper.Reply(ctx, b, msg, v.render(p, f)) diff --git a/internal/modules/thoitiet/thoitiet_test.go b/internal/modules/weather/weather_test.go similarity index 89% rename from internal/modules/thoitiet/thoitiet_test.go rename to internal/modules/weather/weather_test.go index 87a1be0..fe80185 100644 --- a/internal/modules/thoitiet/thoitiet_test.go +++ b/internal/modules/weather/weather_test.go @@ -1,4 +1,4 @@ -package thoitiet +package weather import ( "context" @@ -82,12 +82,12 @@ func stubOpenMeteo(t *testing.T, f *fakeOpenMeteo) { t.Cleanup(func() { geocodeURL, forecastURL = origGeocode, origForecast }) } -func installThoitiet(t *testing.T) *testutil.RecordingBot { +func installWeather(t *testing.T) *testutil.RecordingBot { t.Helper() rb := testutil.NewRecordingBot(t) - mod := New(modules.Deps{Store: storage.NewMemoryProvider().Collection("thoitiet")}) + mod := New(modules.Deps{Store: storage.NewMemoryProvider().Collection(CollectionName)}) reg := &modules.Registry{ - Modules: []modules.Module{{Name: "thoitiet", Commands: mod.Commands}}, + Modules: []modules.Module{{Name: CollectionName, Commands: mod.Commands}}, AllCommands: map[string]modules.Command{}, } for _, c := range mod.Commands { @@ -103,28 +103,39 @@ func send(rb *testutil.RecordingBot, text string) string { } func TestCommands_RegistrationAndParameters(t *testing.T) { - mod := New(modules.Deps{}) - want := []string{"thoitiethomnay", "thoitiet", "thoitietngaymai", "thoitiettuannay"} + mod := New(modules.Deps{Store: storage.NewMemoryProvider().Collection(CollectionName)}) + want := []struct{ name, parameters string }{ + {"thoitiethomnay", "[location...]"}, + {"thoitiet", "[location...]"}, + {"thoitietngaymai", "[location...]"}, + {"thoitiettuannay", "[location...]"}, + {"thuyvan", ""}, + {"thuyvan_subscribe", ""}, + {"thuyvan_unsubscribe", ""}, + } if len(mod.Commands) != len(want) { t.Fatalf("commands = %d, want %d", len(mod.Commands), len(want)) } for i, c := range mod.Commands { - if c.Name != want[i] { - t.Errorf("commands[%d] = %q, want %q", i, c.Name, want[i]) + if c.Name != want[i].name { + t.Errorf("commands[%d] = %q, want %q", i, c.Name, want[i].name) } - if c.Parameters != "[location...]" { - t.Errorf("/%s parameters = %q", c.Name, c.Parameters) + if c.Parameters != want[i].parameters { + t.Errorf("/%s parameters = %q, want %q", c.Name, c.Parameters, want[i].parameters) } if c.Visibility != modules.VisibilityPublic { t.Errorf("/%s is not public", c.Name) } } + if len(mod.Crons) != 2 || mod.Crons[0].Schedule != "30 3 * * *" || mod.Crons[1].Schedule != "30 5 * * *" { + t.Errorf("crons = %+v, want the 10:30 ICT flood push and its 12:30 retry", mod.Crons) + } } func TestToday_DefaultsToHCMWithoutGeocoding(t *testing.T) { f := &fakeOpenMeteo{forecastBody: forecastFixture} stubOpenMeteo(t, f) - rb := installThoitiet(t) + rb := installWeather(t) want := "🌦️ Thời tiết hôm nay 01/10 — Thành phố Hồ Chí Minh\n" + "Hiện tại: 29°C (cảm giác 36°C), Nhiều mây ☁️\n" + @@ -150,7 +161,7 @@ func TestToday_DefaultsToHCMWithoutGeocoding(t *testing.T) { func TestHourly_ListsNextSixHours(t *testing.T) { f := &fakeOpenMeteo{forecastBody: forecastFixture} stubOpenMeteo(t, f) - rb := installThoitiet(t) + rb := installWeather(t) want := "🕐 Thời tiết 6 giờ tới — Thành phố Hồ Chí Minh\n" + "Hiện tại 10:30: 29°C (cảm giác 36°C), Nhiều mây ☁️\n" + @@ -175,7 +186,7 @@ func TestHourly_ListsNextSixHours(t *testing.T) { func TestHourly_NoUpcomingHoursRepliesError(t *testing.T) { stale := strings.Replace(forecastFixture, `"time":"2026-10-01T10:30"`, `"time":"2026-10-01T16:30"`, 1) stubOpenMeteo(t, &fakeOpenMeteo{forecastBody: stale}) - rb := installThoitiet(t) + rb := installWeather(t) if got := send(rb, "/thoitiet"); got != fetchErrorText { t.Errorf("reply = %q, want %q", got, fetchErrorText) @@ -185,7 +196,7 @@ func TestHourly_NoUpcomingHoursRepliesError(t *testing.T) { func TestTomorrow_GeocodesDiacriticLocation(t *testing.T) { f := &fakeOpenMeteo{geocodeBody: daLatGeocodeFixture, forecastBody: forecastFixture} stubOpenMeteo(t, f) - rb := installThoitiet(t) + rb := installWeather(t) got := send(rb, "/thoitietngaymai Đà Lạt ") want := "🌦️ Thời tiết ngày mai T6 02/10 — Ðà Lạt, Lam Dong\n" + @@ -207,7 +218,7 @@ func TestTomorrow_GeocodesDiacriticLocation(t *testing.T) { func TestWeek_ListsSevenDays(t *testing.T) { f := &fakeOpenMeteo{forecastBody: forecastFixture} stubOpenMeteo(t, f) - rb := installThoitiet(t) + rb := installWeather(t) want := "📅 Thời tiết 7 ngày tới — Thành phố Hồ Chí Minh\n" + "T5 01/10: 24–33°C 🌦️ Mưa rào nhẹ, mưa 70%\n" + @@ -231,7 +242,7 @@ func TestAliasExpansionAndForeignPlace(t *testing.T) { forecastBody: forecastFixture, } stubOpenMeteo(t, f) - rb := installThoitiet(t) + rb := installWeather(t) got := send(rb, "/thoitiettuannay Tokyo") if !strings.HasPrefix(got, "📅 Thời tiết 7 ngày tới — Tokyo, Nhật Bản\n") { @@ -246,7 +257,7 @@ func TestAliasExpansionAndForeignPlace(t *testing.T) { func TestUnknownLocation(t *testing.T) { f := &fakeOpenMeteo{geocodeBody: `{"generationtime_ms":0.3}`} stubOpenMeteo(t, f) - rb := installThoitiet(t) + rb := installWeather(t) if got, want := send(rb, "/thoitiet xyzzy"), `Không tìm thấy địa điểm "xyzzy".`; got != want { t.Errorf("reply = %q, want %q", got, want) @@ -266,7 +277,7 @@ func TestUpstreamFailureRepliesError(t *testing.T) { for name, f := range cases { t.Run(name, func(t *testing.T) { stubOpenMeteo(t, f) - rb := installThoitiet(t) + rb := installWeather(t) cmd := "/thoitiet" if f.geocodeBody != "" { cmd = "/thoitiet Hue" @@ -284,7 +295,7 @@ func TestTomorrow_NeedsTwoDailyRows(t *testing.T) { ).Replace(forecastFixture) f := &fakeOpenMeteo{forecastBody: oneDay} stubOpenMeteo(t, f) - rb := installThoitiet(t) + rb := installWeather(t) if got := send(rb, "/thoitietngaymai"); got != fetchErrorText { t.Errorf("reply = %q, want %q", got, fetchErrorText)