Skip to content

naneos.iotweb

NaneosUploadThread

Bases: Thread

Source code in src/naneos/iotweb/naneos_upload_thread.py
 18
 19
 20
 21
 22
 23
 24
 25
 26
 27
 28
 29
 30
 31
 32
 33
 34
 35
 36
 37
 38
 39
 40
 41
 42
 43
 44
 45
 46
 47
 48
 49
 50
 51
 52
 53
 54
 55
 56
 57
 58
 59
 60
 61
 62
 63
 64
 65
 66
 67
 68
 69
 70
 71
 72
 73
 74
 75
 76
 77
 78
 79
 80
 81
 82
 83
 84
 85
 86
 87
 88
 89
 90
 91
 92
 93
 94
 95
 96
 97
 98
 99
100
class NaneosUploadThread(Thread):
    URL: ClassVar[str] = "https://hg3zkburji.execute-api.eu-central-1.amazonaws.com/prod/proto/v1"
    HEADERS: ClassVar[dict] = {
        "Content-Type": "application/json",
        "Accept": "application/json",
    }

    def __init__(
        self,
        data: dict[int, pd.DataFrame],
        callback: Callable[[bool], None] | None,
    ) -> None:
        """Adding the data that should be uploaded to the database.

        Args:
            data (dict[int, pd.DataFrame]): Data to upload, keyed by device serial number.
            callback (Callable[[bool], None] | None): Called with the upload result.
        """
        super().__init__()
        self.data = data
        self._callback = callback

    def run(self) -> None:
        try:
            ret = self.upload(self.data)

            if self._callback:
                if ret.status_code == 200:
                    self._callback(True)
                else:
                    self._callback(False)
        except Exception as e:
            logger.exception(f"Error in upload: {e}")
            if self._callback:
                self._callback(False)

    @staticmethod
    def get_body(upload_string: str) -> str:
        """The JSON envelope the backend expects around the base64 protobuf."""
        return json.dumps(
            {
                "gateway": "python_webhook",
                "data": upload_string,
                "published_at": datetime.datetime.now(datetime.UTC).isoformat(),
            }
        )

    @staticmethod
    def to_upload_frame(df: pd.DataFrame) -> pd.DataFrame:
        """Prepare one device frame for the backend: index in whole seconds, no inf.

        Frames are indexed by unix time in milliseconds (see naneos.frames).
        The index is rounded, not truncated: the devices sample at ~1Hz with a
        phase of their own, so truncating puts the two samples that straddle a
        second boundary into the same second, where one of them wins, and
        leaves the neighbouring second without a row at all.
        """
        df = df.replace([float("inf"), -float("inf")], 0)
        df.index = pd.Index(
            np.rint(df.index.to_numpy(dtype="float64") / 1e3).astype("int64"),
            name=df.index.name,
        )
        return df

    @classmethod
    def build_combined_entry(cls, data: dict[int, pd.DataFrame], abs_time: int):
        """The protobuf message for a snapshot, with timestamps relative to abs_time."""
        devices = [
            create_proto_device(sn, abs_time, cls.to_upload_frame(df)) for sn, df in data.items()
        ]
        return create_combined_entry(devices=devices, abs_timestamp=abs_time)

    @classmethod
    def upload(cls, data: dict[int, pd.DataFrame]) -> requests.Response:
        abs_time = int(datetime.datetime.now().timestamp())
        combined_entry = cls.build_combined_entry(data, abs_time)

        proto_str = combined_entry.SerializeToString()
        proto_str_base64 = base64.b64encode(proto_str).decode()

        body = cls.get_body(proto_str_base64)
        r = requests.post(cls.URL, headers=cls.HEADERS, data=body, timeout=10)
        return r

__init__(data, callback)

Adding the data that should be uploaded to the database.

Parameters:

Name Type Description Default
data dict[int, DataFrame]

Data to upload, keyed by device serial number.

required
callback Callable[[bool], None] | None

Called with the upload result.

required
Source code in src/naneos/iotweb/naneos_upload_thread.py
25
26
27
28
29
30
31
32
33
34
35
36
37
38
def __init__(
    self,
    data: dict[int, pd.DataFrame],
    callback: Callable[[bool], None] | None,
) -> None:
    """Adding the data that should be uploaded to the database.

    Args:
        data (dict[int, pd.DataFrame]): Data to upload, keyed by device serial number.
        callback (Callable[[bool], None] | None): Called with the upload result.
    """
    super().__init__()
    self.data = data
    self._callback = callback

build_combined_entry(data, abs_time) classmethod

The protobuf message for a snapshot, with timestamps relative to abs_time.

Source code in src/naneos/iotweb/naneos_upload_thread.py
82
83
84
85
86
87
88
@classmethod
def build_combined_entry(cls, data: dict[int, pd.DataFrame], abs_time: int):
    """The protobuf message for a snapshot, with timestamps relative to abs_time."""
    devices = [
        create_proto_device(sn, abs_time, cls.to_upload_frame(df)) for sn, df in data.items()
    ]
    return create_combined_entry(devices=devices, abs_timestamp=abs_time)

get_body(upload_string) staticmethod

The JSON envelope the backend expects around the base64 protobuf.

Source code in src/naneos/iotweb/naneos_upload_thread.py
54
55
56
57
58
59
60
61
62
63
@staticmethod
def get_body(upload_string: str) -> str:
    """The JSON envelope the backend expects around the base64 protobuf."""
    return json.dumps(
        {
            "gateway": "python_webhook",
            "data": upload_string,
            "published_at": datetime.datetime.now(datetime.UTC).isoformat(),
        }
    )

to_upload_frame(df) staticmethod

Prepare one device frame for the backend: index in whole seconds, no inf.

Frames are indexed by unix time in milliseconds (see naneos.frames). The index is rounded, not truncated: the devices sample at ~1Hz with a phase of their own, so truncating puts the two samples that straddle a second boundary into the same second, where one of them wins, and leaves the neighbouring second without a row at all.

Source code in src/naneos/iotweb/naneos_upload_thread.py
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
@staticmethod
def to_upload_frame(df: pd.DataFrame) -> pd.DataFrame:
    """Prepare one device frame for the backend: index in whole seconds, no inf.

    Frames are indexed by unix time in milliseconds (see naneos.frames).
    The index is rounded, not truncated: the devices sample at ~1Hz with a
    phase of their own, so truncating puts the two samples that straddle a
    second boundary into the same second, where one of them wins, and
    leaves the neighbouring second without a row at all.
    """
    df = df.replace([float("inf"), -float("inf")], 0)
    df.index = pd.Index(
        np.rint(df.index.to_numpy(dtype="float64") / 1e3).astype("int64"),
        name=df.index.name,
    )
    return df

download_from_iotweb(name, serial_number, start, stop, token)

Download your data from influxdb.naneos.ch. 1 Month of data takes about 30 seconds to download and uses about 100 MB of data.

You need to have a token to access the data. Ask mario.huegi@naneos.ch for your read token. We kindly ask you to not overuse our server. If you need to download the same data in a recuring pattern, contact us.

Parameters:

Name Type Description Default
name str

Name of the influx bucket.

required
serial_number str

Serial number of your device as string.

required
start datetime

Start date of the data you want to download.

required
stop datetime

End date of the data you want to download.

required
token str

Your read token. Do not push your token to public repositories.

required

Returns:

Type Description
DataFrame

pd.DataFrame: Dataframe with your data.

Source code in src/naneos/iotweb/download/downloader.py
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
def download_from_iotweb(
    name: str, serial_number: str, start: dt.datetime, stop: dt.datetime, token: str
) -> pd.DataFrame:
    """Download your data from influxdb.naneos.ch.
    1 Month of data takes about 30 seconds to download and uses about 100 MB of data.

    You need to have a token to access the data.
    Ask mario.huegi@naneos.ch for your read token.
    We kindly ask you to not overuse our server.
    If you need to download the same data in a recuring pattern, contact us.

    Args:
        name (str): Name of the influx bucket.
        serial_number (str): Serial number of your device as string.
        start (dt.datetime): Start date of the data you want to download.
        stop (dt.datetime): End date of the data you want to download.
        token (str): Your read token. Do not push your token to public repositories.

    Returns:
        pd.DataFrame: Dataframe with your data.
    """
    timestamps = create_start_stop_timestamp(start, stop)

    dfs = []

    with InfluxDBClient(url=URL_INFLUX, org=ORG_INFLUX, token=token) as client:
        for t1, t2 in timestamps:
            query = get_query(name, serial_number, t1, t2)

            df = client.query_api().query_data_frame(query)

            if isinstance(df, list):
                dfs.extend(df)
            elif isinstance(df, pd.DataFrame):
                dfs.append(df)
            else:
                logger.warning(f"Unknown type: {type(df)}")

    df = pd.concat(dfs, axis=0)
    df.set_index("_time", inplace=True)
    df.drop(["result", "table"], axis=1, inplace=True)

    return df