Skip to content

naneos.iotweb.naneos_upload_thread

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