diff --git a/src/main/scala/Main.scala b/src/main/scala/Main.scala index 589b31a..9b5384a 100644 --- a/src/main/scala/Main.scala +++ b/src/main/scala/Main.scala @@ -37,6 +37,16 @@ object Main extends IOApp { } yield () fetchStationsTask = scheduler.scheduleTask("Fetch Weather Stations", List(11,13,23,30), fetchWeatherStations).compile.drain + // Real, free, keyless alternative to the private FTP feed above (see + // docs/UPDATE_ROADMAP.md). Runs alongside it for now rather than + // replacing it, so the two can be compared before the FTP path is removed. + openDataStationService <- lvgmc.OpenDataStationService.of + fetchOpenDataStations = for { + stationData <- openDataStationService.fetchStationObservations + _ <- postgresService.insertInWeatherTable(stationData) + } yield () + fetchOpenDataStationsTask = scheduler.scheduleTask("Fetch Open Data Stations", List(15,45), fetchOpenDataStations).compile.drain + cleanupTask = scheduler.scheduleTask("Cleanup old Grib", List(41), dataService.deleteOldForecasts()).compile.drain fetchGrib <- dmi.FetchService.of(dataService) @@ -49,7 +59,7 @@ object Main extends IOApp { if (sys.env.get("ENABLE_SCHEDULED_JOBS").exists(_.equalsIgnoreCase("false"))) IO.never[Unit] else - (fetchStationsTask, cleanupTask, fetchGribTask).parMapN((_, _, _) => ()) + (fetchStationsTask, fetchOpenDataStationsTask, cleanupTask, fetchGribTask).parMapN((_, _, _, _) => ()) exitCode <- (serverTask, scheduledTasks).parMapN((_, _) => ExitCode.Success) } yield exitCode diff --git a/src/main/scala/fetch/lvgmc/OpenDataStationService.scala b/src/main/scala/fetch/lvgmc/OpenDataStationService.scala new file mode 100644 index 0000000..b00f629 --- /dev/null +++ b/src/main/scala/fetch/lvgmc/OpenDataStationService.scala @@ -0,0 +1,145 @@ +package fetch.lvgmc + +import cats.effect.IO +import cats.syntax.all._ +import io.circe.Json +import io.circe.parser +import org.http4s.Uri +import org.http4s.client.Client +import org.http4s.ember.client.EmberClientBuilder +import org.typelevel.log4cats.Logger +import org.typelevel.log4cats.slf4j.Slf4jLogger +import parse.csv.{WeatherData, WeatherStationData} + +import java.time.LocalDateTime +import scala.util.Try + +object OpenDataStationService { + def of: IO[OpenDataStationService] = + Slf4jLogger.create[IO].map(logger => new OpenDataStationService(logger)) +} + +// Fetches current LVĢMC station observations from Latvia's open data portal +// (data.gov.lv, dataset "hidrometeorologiskie-noverojumi") rather than the +// private LVGMC FTP feed used by fetch.lvgmc.FetchService. Free, keyless, +// same underlying LVĢMC data. See docs/UPDATE_ROADMAP.md for the phased plan +// this is part of. +final class OpenDataStationService private (logger: Logger[IO]) { + private val pageSize = 5000 + + private val observationsUrl = sys.env.getOrElse( + "LVGMC_OPENDATA_OBSERVATIONS_URL", + "https://data.gov.lv/dati/api/3/action/datastore_search?resource_id=17460efb-ae99-4d1d-8144-1068f184b05f" + ) + + // Matched against LVĢMC's open-data parameter dictionary (resource + // 38b462ac-08b9-4168-9d6e-cbaedc2e775d) using the hourly vs + // observation-time semantics documented for each WeatherData field. + // visibilityMin, dewPoint, and sunDuration have no equivalent in this + // dataset and are left as None below. + private val fieldAbbreviations = + List("HATMX", "HATMN", "HTDRY", "HPRAB", "WNS10", "HWSMX", "VSBAV", "HSNOW", "PRSL", "RLH") + + // Verified 2026-08-23 against the open station-metadata resource + // (c32c7afd-0d05-44fd-8b24-1de85b4bf11d): 32 of the app's tracked cities + // have a direct match. RIGASLU ("Rīga Universitāte") is the representative + // Rīga station. SELIEPA ("Liepāja piekraste") is a duplicate coastal + // station and is intentionally not mapped. Cēsis, Jēkabpils, Talsi, and + // Valmiera have no station in this dataset; Jēkabpils/Talsi already use + // Zīlāni/Stende as substitutes elsewhere (see docs/PRODUCT_WORKFLOWS.md). + private val stationCities: Map[String, String] = Map( + "KALNCIEM" -> "Kalnciems", "SIGULDA" -> "Sigulda", "RIVE99PA" -> "Ventspils", + "RIJE99PA" -> "Jelgava", "RIDAGDA" -> "Dagda", "PIEDRUJA" -> "Piedruja", + "VICAKI" -> "Vičaki", "DAUGAVGR" -> "Daugavgrīva", "RIZI99PA" -> "Zīlāni", + "RIME99MS" -> "Mērsrags", "RIAI99PA" -> "Ainaži", "RIGASLU" -> "Rīga", + "RIKO99PA" -> "Kolka", "RIREZEKN" -> "Rēzekne", "RIBA99PA" -> "Bauska", + "RIDO99MS" -> "Dobele", "RIST99PA" -> "Stende", "RIDM99MS" -> "Daugavpils", + "RISA99PA" -> "Saldus", "RIZO99MS" -> "Zosēni", "RIRU99PA" -> "Rūjiena", + "KULDIGA" -> "Kuldīga", "RIPR99PA" -> "Priekuļi", "RIAL99MS" -> "Alūksne", + "RIGU99MS" -> "Gulbene", "RILP99PA" -> "Liepāja", "SILI" -> "Sīļi", + "LIELPECI" -> "Lielpēči", "RIMADONA" -> "Madona", "RIPA99PA" -> "Pāvilosta", + "RUCAVA" -> "Rucava", "RISE99MS" -> "Skulte", "RISI99PA" -> "Skrīveri", + ) + + def fetchStationObservations: IO[List[WeatherStationData]] = + EmberClientBuilder.default[IO].build.use { client => + for { + baseUri <- Uri.fromString(observationsUrl).liftTo[IO] + rows <- fetchRows(client, baseUri) + stationData = buildStationData(rows) + _ <- logger.info(s"Loaded ${stationData.size} open-data station observation rows") + } yield stationData + } + + private def fetchRows( + client: Client[IO], + baseUri: Uri, + offset: Int = 0, + accumulated: Vector[Json] = Vector.empty + ): IO[Vector[Json]] = { + val filters = Json.obj("ABBREVIATION" -> Json.arr(fieldAbbreviations.map(Json.fromString): _*)).noSpaces + val uri = baseUri + .withQueryParam("limit", pageSize.toString) + .withQueryParam("offset", offset.toString) + .withQueryParam("filters", filters) + + client.expect[String](uri).flatMap(parsePage).flatMap { case (total, rows) => + val allRows = accumulated ++ rows + if (allRows.size >= total || rows.isEmpty) IO.pure(allRows) + else fetchRows(client, baseUri, offset + rows.size, allRows) + } + } + + private def parsePage(body: String): IO[(Int, Vector[Json])] = + IO.fromEither( + parser.parse(body) + .leftMap(error => new RuntimeException(s"Unable to parse open-data observations: ${error.message}")) + .flatMap { json => + val result = json.hcursor.downField("result") + (for { + total <- result.downField("total").as[Int] + records <- result.downField("records").as[Vector[Json]] + } yield total -> records).leftMap(error => new RuntimeException(s"Invalid open-data observations response: ${error.message}")) + } + ) + + private case class ObservationCell(city: String, timestamp: LocalDateTime, abbreviation: String, value: Double) + + private def buildStationData(rows: Vector[Json]): List[WeatherStationData] = { + val cells = rows.flatMap { row => + val cursor = row.hcursor + for { + stationId <- cursor.downField("STATION_ID").as[String].toOption + city <- stationCities.get(stationId) + abbreviation <- cursor.downField("ABBREVIATION").as[String].toOption + datetimeStr <- cursor.downField("DATETIME").as[String].toOption + timestamp <- Try(LocalDateTime.parse(datetimeStr)).toOption + value <- cursor.downField("VALUE").as[Double].toOption + } yield ObservationCell(city, timestamp, abbreviation, value) + } + + cells + .groupBy(cell => (cell.city, cell.timestamp)) + .map { case ((city, timestamp), group) => + val values = group.map(cell => cell.abbreviation -> cell.value).toMap + val weather = WeatherData( + tempMax = values.get("HATMX"), + tempMin = values.get("HATMN"), + tempAvg = values.get("HTDRY"), + precipitation = values.get("HPRAB"), + windAvg = values.get("WNS10"), + windMax = values.get("HWSMX"), + visibilityMin = None, + visibilityAvg = values.get("VSBAV"), + snowAvg = values.get("HSNOW"), + atmPressure = values.get("PRSL"), + dewPoint = None, + humidity = values.get("RLH"), + sunDuration = None, + phenomena = List.empty, + ) + WeatherStationData(city, timestamp, weather) + } + .toList + } +} diff --git a/src/main/scala/fetch/lvgmc/OpenDataStationServiceTest.scala b/src/main/scala/fetch/lvgmc/OpenDataStationServiceTest.scala new file mode 100644 index 0000000..a8a287f --- /dev/null +++ b/src/main/scala/fetch/lvgmc/OpenDataStationServiceTest.scala @@ -0,0 +1,32 @@ +package fetch.lvgmc + +import cats.effect.IO +import cats.effect.unsafe.implicits.global +import cats.syntax.all._ +import db.{DBConnection, PostgresService} + +// Manual verification harness, matching the existing FetchServiceTest/ +// GribParserTest pattern. +// Fetch only: sbt "runMain fetch.lvgmc.OpenDataStationServiceTest" +// Fetch and write: sbt "runMain fetch.lvgmc.OpenDataStationServiceTest --write" +object OpenDataStationServiceTest { + def main(args: Array[String]): Unit = { + val program = for { + service <- OpenDataStationService.of + stationData <- service.fetchStationObservations + _ <- IO.println(s"Fetched ${stationData.size} (city, hour) rows") + _ <- IO.println(s"Distinct cities: ${stationData.map(_.city).distinct.sorted.mkString(", ")}") + _ <- stationData.sortBy(_.timestamp).reverse.take(5).traverse_(row => IO.println(row)) + _ <- if (args.contains("--write")) { + for { + transactor <- DBConnection.transactor[IO] + postgresService <- PostgresService.of(transactor) + rowCount <- postgresService.insertInWeatherTable(stationData) + _ <- IO.println(s"Wrote $rowCount rows to the weather table") + } yield () + } else IO.unit + } yield () + + program.unsafeRunSync() + } +}