Metadata-Version: 2.5
Name: omixom-data
Version: 0.3.1
Summary: Cliente Python de la Omixom Data API v3: lectura de series y réplica incremental de mediciones.
Author: Omixom
License: MIT
Requires-Python: >=3.11
Requires-Dist: httpx>=0.27
Requires-Dist: pydantic>=2.7
Description-Content-Type: text/markdown

# omixom-data

Cliente Python de la [Omixom Data API v3](https://clima.omixom.com/api/v3/docs) para la
consulta de series de mediciones y la réplica incremental que mantiene una copia local
sincronizada.

```bash
pip install omixom-data
```

Requiere Python 3.11 o superior y un token de acceso de la red Omixom.

## Réplica incremental

El caso de uso principal es replicar las mediciones de un conjunto de estaciones y mantener esa
réplica actualizada en el tiempo. El `Feed` encapsula el protocolo completo (bootstrap paginado,
cursores, coherencia del snapshot) y entrega el resultado en batches, donde cada batch
contiene una página de eventos junto con el estado que queda tras aplicarla.

```python
from omixom_data import Client, MeasurementDeleted

with Client(token="...") as client:
    feed = client.feed([30125, 30126])  # la historia completa de cada estación

    for batch in feed.batches():  # bootstrap y cambios pendientes
        for event in batch.events:
            if isinstance(event, MeasurementDeleted):
                store.delete(event.station, event.module, event.time)
            else:
                store.upsert(event.station, event.module, event.time, event.value)
        save(batch.state.model_dump_json())  # checkpoint alineado a la página
```

El contrato es aplicar primero y persistir el estado después, idealmente ambos dentro de una
transacción del almacenamiento propio. Con ese orden, una interrupción en cualquier punto
(incluso durante un backfill extenso) tiene como peor caso una página re-aplicada, y aplicar
eventos es idempotente (`MeasurementUpserted` fija el valor, `MeasurementDeleted` lo elimina y
re-aplicarlos no altera el resultado). Cada batch persistido consolida el avance; interrumpir el
recorrido entre batches es siempre seguro.

## Sincronización continua

Un proceso de sincronización continua reúne las piezas anteriores, retomando el estado
persistido si existe (`client.resume`), procesando los batches pendientes y esperando el
intervalo de consulta.
El estado se guarda serializado en el mismo store que recibe los datos, y cada batch se
confirma en una única transacción que abarca los eventos y el checkpoint.

```python
import time

from omixom_data import Client, FeedState, MeasurementDeleted, MeasurementUpserted

POLL_INTERVAL = 300  # segundos
STATIONS = [30125, 30126]


def apply(event: MeasurementUpserted | MeasurementDeleted) -> None:
    if isinstance(event, MeasurementDeleted):
        store.delete(event.station, event.module, event.time)
    else:
        store.upsert(event.station, event.module, event.time, event.value)


with Client(token="...") as client:
    saved = store.load_state()  # el JSON guardado en la corrida anterior, o None
    if saved is not None:
        feed = client.resume(FeedState.model_validate_json(saved))
    else:
        feed = client.feed(STATIONS)

    while True:
        for batch in feed.batches():  # solo lo pendiente desde el último checkpoint
            with store.transaction():  # eventos y checkpoint juntos, todo o nada
                for event in batch.events:
                    apply(event)
                store.save_state(batch.state.model_dump_json())
        time.sleep(POLL_INTERVAL)
```

Con esa transacción por batch el checkpoint queda exactamente alineado con lo aplicado,
porque una falla a mitad de un batch revierte todo y la corrida siguiente retoma desde el checkpoint
anterior, sin estados intermedios. Si el estado no puede vivir en el store y se guarda aparte
(por ejemplo, en un archivo JSON local, escrito de forma atómica vía un temporal y
`Path.replace`), la garantía se relaja, y una interrupción entre aplicar y guardar implica
re-aplicar la última página en la corrida siguiente, lo que no altera el resultado por la
idempotencia de los eventos.

`batches()` puede invocarse repetidamente sobre el mismo feed, ya que cada llamada avanza
hasta alcanzar el estado actual y la siguiente retoma desde ese punto, por lo que el ciclo del
proceso es invocarlo y esperar el intervalo, manteniendo el `Client` abierto para reutilizar
las conexiones. El avance se registra por módulo (el estado guarda un cursor por serie) y los
requests se agrupan por estación, con el bootstrap de varias estaciones avanzando round-robin
para distribuir los límites de uso.

## Acotar el rango

```python
from datetime import UTC, datetime

feed = client.feed([30125], date_from=datetime(2023, 1, 1, tzinfo=UTC))  # desde 2023
feed = client.feed(
    [30125],
    date_from=datetime(2023, 1, 1, tzinfo=UTC),
    date_to=datetime(2024, 1, 1, tzinfo=UTC),  # 2023 completo
)
```

Sin fechas, cada estación se replica desde su fecha de instalación. Con `date_from` únicamente,
la réplica continúa recibiendo los datos nuevos; con ambas fechas la ventana queda cerrada,
aunque las correcciones sobre datos de esa ventana siguen llegando. `modules` restringe la
réplica a esos módulos, de cualquiera de las estaciones indicadas (una estación sin módulos
seleccionados no genera requests); el conjunto replicado queda fijado al crear el feed, por lo
que un sensor instalado con posterioridad no se incorpora automáticamente (ver la sección
siguiente). `categories` filtra el tipo de dato a replicar; con ese filtro, un punto corregido
hacia una categoría no seleccionada se recibe como borrado, dado que salió de la vista
replicada.

## Incorporar módulos a un feed existente

Un módulo instalado después de crear el feed se incorpora con `add_modules`, sobre un feed nuevo
o retomado; sin lista de módulos se incorporan todos los de la estación que aún no estén
rastreados.

```python
feed = client.resume(FeedState.model_validate_json(saved))
feed.add_modules(30125)  # o add_modules(30125, [4812]) para un módulo puntual

for batch in feed.batches():  # el módulo nuevo hace su bootstrap; el resto continúa
    apply_all(batch.events)
    save(batch.state.model_dump_json())
```

El método devuelve los ids incorporados (los ya rastreados se omiten) y levanta `ValueError`
ante un módulo que no pertenece a la estación. Sirve también para incorporar una estación nueva
al feed. El alta queda persistida con el `state` del batch siguiente; si el proceso se
interrumpe antes, la corrida siguiente debe repetir la llamada.

## Lecturas puntuales

Las lecturas puntuales resuelven consultas únicas, sin réplica de por medio.

```python
client.stations()  # estaciones accesibles con el token
client.station(30125)  # fecha de instalación y módulos

for point in client.series(  # la serie completa de una ventana, en streaming
    30125,
    date_from=datetime(2024, 1, 1, tzinfo=UTC),
    date_to=datetime(2024, 2, 1, tzinfo=UTC),
):
    print(point.module, point.time, point.value)
```

`series` recorre internamente todas las páginas repitiendo el cursor de la primera, de modo que
el resultado completo es un corte coherente de la base aun ante escrituras concurrentes. Para
materializar la serie completa, `list(client.series(...))`.

Por debajo existe el acceso crudo página por página (`series_page`, `changes_page`), donde el
manejo del cursor queda a cargo del caller. Para paginar de forma coherente se repite el
request con `date_from` igual al `next_from` recibido y `cursor` igual al `cursor` recibido.

## Errores y límites de uso

Los errores de la API se levantan como excepciones tipadas bajo `OmixomDataError`
(`AuthenticationError`, `NotFoundError`, `InvalidRequestError`, `RateLimitedError` y
`ServerError`). Ante un 429 el cliente espera el tiempo indicado en `Retry-After` y reintenta
automáticamente; `Client(..., wait_on_rate_limit=False)` desactiva esa espera y levanta
`RateLimitedError` con el tiempo sugerido en `retry_after`.
