• Home
  • Features
  • Pricing
  • Docs
  • Announcements
  • Sign In

kimata / my-py-lib / 20586913338

30 Dec 2025 01:53AM UTC coverage: 65.787% (-0.05%) from 65.835%
20586913338

push

github

kimata
feat(sensor_data): SensorDataResult に診断情報フィールドを追加

- raw_record_count: 取得した生レコード数
- null_count: None だったレコード数
- error_message: エラー発生時のメッセージ
- get_diagnostic_message(): 診断メッセージを生成するメソッド

また、データ切り詰め処理を改善し、データ不足時に全データが
消えないように修正。

🤖 Generated with [Claude Code](https://claude.com/claude-code)

Co-Authored-By: Claude Opus 4.5 <noreply@anthropic.com>

16 of 27 new or added lines in 1 file covered. (59.26%)

2915 of 4431 relevant lines covered (65.79%)

0.66 hits per line

Source File
Press 'n' to go to next uncovered line, 'b' for previous

66.95
/src/my_lib/sensor_data.py
1
#!/usr/bin/env python3
2
"""
3
InfluxDB からデータを取得します。
4

5
Usage:
6
  sensor_data.py [-c CONFIG] [-m MODE] [-i DB_SPEC] [-s SENSOR_SPEC] [-f FIELD] [-e EVERY] [-w WINDOW]
7
                 [-p HOURS] [-D]
8

9
Options:
10
  -c CONFIG         : CONFIG を設定ファイルとして読み込んで実行します。[default: tests/data/config.example.yaml]
11
  -m MODE           : データ取得モード。(data, day_sum, hour_sum, minute_sum のいずれか) [default: data]
12
  -i DB_SPEC        : 設定ファイルの中で InfluxDB の設定が書かれているパス。[default: sensor.influxdb]
13
  -s SENSOR_SPEC    : 設定ファイルの中で取得対象のデータの設定が書かれているパス。[default: sensor.lux]
14
  -f FIELD          : 取得するフィールド。[default: lux]
15
  -e EVERY          : 何分ごとのデータを取得するか。[default: 1]
16
  -w WINDOWE        : 算出に使うウィンドウ。[default: 5]
17
  -p PERIOD         : 積算(sum)モードの場合に、過去どのくらいの分を取得するか。[default: 1]
18
  -D                : デバッグモードで動作します。
19
"""
20

21
from __future__ import annotations
1✔
22

23
import asyncio
1✔
24
import datetime
1✔
25
import logging
1✔
26
import os
1✔
27
import time
1✔
28
from dataclasses import dataclass, field
1✔
29
from typing import Any, TypedDict, cast
1✔
30

31

32
class InfluxDBConfig(TypedDict):
1✔
33
    """InfluxDB 接続設定"""
34

35
    url: str
1✔
36
    token: str
1✔
37
    org: str
1✔
38
    bucket: str
1✔
39

40

41
import influxdb_client
1✔
42
from influxdb_client.client.flux_table import TableList
1✔
43

44
import my_lib.time
1✔
45

46

47
@dataclass(frozen=True)
1✔
48
class SensorDataResult:
1✔
49
    """センサーデータ取得結果
50

51
    Attributes:
52
        value: センサー値のリスト
53
        time: タイムスタンプのリスト
54
        valid: データが有効かどうか
55
        raw_record_count: 取得した生レコード数(処理前)
56
        null_count: None だったレコード数
57
        error_message: エラー発生時のメッセージ
58
    """
59

60
    value: list[float] = field(default_factory=list)
1✔
61
    time: list[datetime.datetime] = field(default_factory=list)
1✔
62
    valid: bool = False
1✔
63
    raw_record_count: int = 0
1✔
64
    null_count: int = 0
1✔
65
    error_message: str | None = None
1✔
66

67
    def get_diagnostic_message(self) -> str:
1✔
68
        """診断メッセージを生成"""
NEW
69
        if self.error_message:
×
NEW
70
            return f"接続エラー: {self.error_message}"
×
NEW
71
        if self.raw_record_count == 0:
×
NEW
72
            return "データなし: クエリ結果が空でした"
×
NEW
73
        if self.null_count == self.raw_record_count:
×
NEW
74
            return f"全データがNone: {self.raw_record_count}件すべてがNoneでした"
×
NEW
75
        if self.null_count > 0:
×
NEW
76
            return (
×
77
                f"一部データがNone: {self.raw_record_count}件中{self.null_count}件がNone、"
78
                f"有効データ{len(self.value)}件"
79
            )
NEW
80
        return f"データ取得成功: {len(self.value)}件"
×
81

82

83
@dataclass(frozen=True)
1✔
84
class DataRequest:
1✔
85
    """センサーデータ取得リクエスト"""
86

87
    measure: str
1✔
88
    hostname: str
1✔
89
    field: str
1✔
90
    start: str = "-30h"
1✔
91
    stop: str = "now()"
1✔
92
    every_min: int = 1
1✔
93
    window_min: int = 3
1✔
94
    create_empty: bool = True
1✔
95
    last: bool = False
1✔
96

97

98
# NOTE: データが欠損している期間も含めてデータを敷き詰めるため、
99
# timedMovingAverage を使う。timedMovingAverage の計算の結果、データが後ろに
100
# ずれるので、あらかじめ offset を使って前にずらしておく。
101
FLUX_QUERY = """
1✔
102
from(bucket: "{bucket}")
103
|> range(start: {start}, stop: {stop})
104
    |> filter(fn:(r) => r._measurement == "{measure}")
105
    |> filter(fn: (r) => r.hostname == "{hostname}")
106
    |> filter(fn: (r) => r["_field"] == "{field}")
107
    |> aggregateWindow(every: {window}m, offset:-{window}m, fn: mean, createEmpty: {create_empty})
108
    |> fill(usePrevious: true)
109
    |> timedMovingAverage(every: {every}m, period: {window}m)
110
"""
111

112
FLUX_QUERY_WITHOUT_AGGREGATION = """
1✔
113
from(bucket: "{bucket}")
114
|> range(start: {start}, stop: {stop})
115
    |> filter(fn:(r) => r._measurement == "{measure}")
116
    |> filter(fn: (r) => r.hostname == "{hostname}")
117
    |> filter(fn: (r) => r["_field"] == "{field}")
118
    |> fill(usePrevious: true)
119
"""
120

121
FLUX_SUM_QUERY = """
1✔
122
from(bucket: "{bucket}")
123
    |> range(start: {start}, stop: {stop})
124
    |> filter(fn:(r) => r._measurement == "{measure}")
125
    |> filter(fn: (r) => r.hostname == "{hostname}")
126
    |> filter(fn: (r) => r["_field"] == "{field}")
127
    |> aggregateWindow(every: {every}m, offset:-{every}m, fn: mean, createEmpty: {create_empty})
128
    |> filter(fn: (r) => exists r._value)
129
    |> fill(usePrevious: true)
130
    |> reduce(
131
        fn: (r, accumulator) => ({{sum: r._value + accumulator.sum, count: accumulator.count + 1}}),
132
        identity: {{sum: 0.0, count: 0}},
133
    )
134
"""
135

136
FLUX_EVENT_QUERY = """
1✔
137
from(bucket: "{bucket}")
138
    |> range(start: {start})
139
    |> filter(fn: (r) => r._measurement == "{measure}")
140
    |> filter(fn: (r) => r.hostname == "{hostname}")
141
    |> filter(fn: (r) => r["_field"] == "{field}")
142
    |> map(fn: (r) => ({{ r with _value: if r._value then 1 else 0 }}))
143
    |> difference()
144
    |> filter(fn: (r) => r._value == 1)
145
    |> sort(columns: ["_time"], desc: true)
146
    |> limit(n: 1)
147
"""
148

149

150
def _process_query_results(
1✔
151
    table_list: list[Any], create_empty: bool, last: bool, every_min: int, window_min: int
152
) -> SensorDataResult:
153
    """共通のクエリ結果処理ロジック"""
154
    data_list = []
1✔
155
    time_list = []
1✔
156
    localtime_offset = datetime.timedelta(hours=9)
1✔
157

158
    # 診断情報
159
    raw_record_count = 0
1✔
160
    null_count = 0
1✔
161

162
    if len(table_list) != 0:
1✔
163
        raw_record_count = len(table_list[0].records)
1✔
164
        for record in table_list[0].records:
1✔
165
            # NOTE: aggregateWindow(createEmpty: true) と fill(usePrevious: true) の組み合わせ
166
            # だとタイミングによって、先頭に None が入る
167
            if record.get_value() is None:
1✔
168
                logging.debug("DELETE %s", record.get_time() + localtime_offset)
1✔
169
                null_count += 1
1✔
170
                continue
1✔
171

172
            data_list.append(record.get_value())
1✔
173
            time_list.append(record.get_time() + localtime_offset)
1✔
174

175
    if create_empty and not last:
1✔
176
        # NOTE: aggregateWindow(createEmpty: true) と timedMovingAverage を使うと、
177
        # 末尾に余分なデータが入るので取り除く
178
        every_min = int(every_min)
1✔
179
        window_min = int(window_min)
1✔
180
        if window_min > every_min:
1✔
181
            trim_count = window_min - every_min
1✔
182
            # データが十分にある場合のみ切り詰め
183
            if len(data_list) > trim_count:
1✔
184
                data_list = data_list[:-trim_count]
1✔
185
                time_list = time_list[:-trim_count]
1✔
186
            else:
187
                logging.warning(
1✔
188
                    "Insufficient data to trim: data_count=%d, trim_count=%d",
189
                    len(data_list),
190
                    trim_count,
191
                )
192

193
    logging.debug("data count = %s", len(time_list))
1✔
194
    return SensorDataResult(
1✔
195
        value=data_list,
196
        time=time_list,
197
        valid=len(time_list) != 0,
198
        raw_record_count=raw_record_count,
199
        null_count=null_count,
200
    )
201

202

203
def _fetch_data_impl(  # noqa: PLR0913
1✔
204
    db_config: InfluxDBConfig,
205
    template: str,
206
    measure: str,
207
    hostname: str,
208
    field: str,
209
    start: str,
210
    stop: str,
211
    every: int,
212
    window: int,
213
    create_empty: bool,
214
    last: bool = False,
215
) -> TableList:
216
    client = None
×
217
    try:
×
218
        token = os.environ.get("INFLUXDB_TOKEN", db_config["token"])
×
219

220
        query = template.format(
×
221
            bucket=db_config["bucket"],
222
            measure=measure,
223
            hostname=hostname,
224
            field=field,
225
            start=start,
226
            stop=stop,
227
            every=every,
228
            window=window,
229
            create_empty=str(create_empty).lower(),
230
        )
231
        if last:
×
232
            query += " |> last()"
×
233

234
        logging.debug("Flux query = %s", query)
×
235
        client = influxdb_client.InfluxDBClient(  # type: ignore[attr-defined]
×
236
            url=db_config["url"], token=token, org=db_config["org"]
237
        )
238
        query_api = client.query_api()
×
239

240
        return query_api.query(query=query)
×
241
    except Exception:
×
242
        logging.exception("Failed to fetch data")
×
243
        raise
×
244
    finally:
245
        if client is not None:
×
246
            client.close()
×
247

248

249
async def _fetch_data_impl_async(  # noqa: PLR0913
1✔
250
    db_config: InfluxDBConfig,
251
    template: str,
252
    measure: str,
253
    hostname: str,
254
    field: str,
255
    start: str,
256
    stop: str,
257
    every: int,
258
    window: int,
259
    create_empty: bool,
260
    last: bool = False,
261
) -> TableList:
262
    """非同期版のデータ取得実装"""
263
    loop = asyncio.get_event_loop()
×
264
    return await loop.run_in_executor(
×
265
        None,
266
        _fetch_data_impl,
267
        db_config,
268
        template,
269
        measure,
270
        hostname,
271
        field,
272
        start,
273
        stop,
274
        every,
275
        window,
276
        create_empty,
277
        last,
278
    )
279

280

281
def fetch_data(  # noqa: PLR0913
1✔
282
    db_config: InfluxDBConfig,
283
    measure: str,
284
    hostname: str,
285
    field: str,
286
    start: str = "-30h",
287
    stop: str = "now()",
288
    every_min: int = 1,
289
    window_min: int = 3,
290
    create_empty: bool = True,
291
    last: bool = False,
292
) -> SensorDataResult:
293
    time_start = time.time()
1✔
294
    logging.debug(
1✔
295
        (
296
            "Fetch data (measure: %s, host: %s, field: %s, "
297
            "start: %s, stop: %s, every: %dmin, window: %dmin, "
298
            "create_empty: %s, last: %s)"
299
        ),
300
        measure,
301
        hostname,
302
        field,
303
        start,
304
        stop,
305
        every_min,
306
        window_min,
307
        create_empty,
308
        last,
309
    )
310

311
    try:
1✔
312
        template = FLUX_QUERY_WITHOUT_AGGREGATION if window_min == 0 else FLUX_QUERY
1✔
313

314
        table_list = _fetch_data_impl(
1✔
315
            db_config,
316
            template,
317
            measure,
318
            hostname,
319
            field,
320
            start,
321
            stop,
322
            every_min,
323
            window_min,
324
            create_empty,
325
            last,
326
        )
327
        time_fetched = time.time()
1✔
328

329
        result = _process_query_results(table_list, create_empty, last, every_min, window_min)
1✔
330

331
        time_finish = time.time()
1✔
332
        if ((time_fetched - time_start) > 1) or ((time_finish - time_fetched) > 0.1):
1✔
333
            logging.warning(
×
334
                "It's taking too long to retrieve the data. (fetch: %.2f sec, modify: %.2f sec)",
335
                time_fetched - time_start,
336
                time_finish - time_fetched,
337
            )
338

339
        return result
1✔
340
    except Exception as e:
1✔
341
        logging.exception("Failed to fetch data")
1✔
342

343
        return SensorDataResult(error_message=str(e))
1✔
344

345

346
async def fetch_data_async(  # noqa: PLR0913
1✔
347
    db_config: InfluxDBConfig,
348
    measure: str,
349
    hostname: str,
350
    field: str,
351
    start: str = "-30h",
352
    stop: str = "now()",
353
    every_min: int = 1,
354
    window_min: int = 3,
355
    create_empty: bool = True,
356
    last: bool = False,
357
) -> SensorDataResult:
358
    """非同期版のfetch_data"""
359
    time_start = time.time()
1✔
360
    logging.debug(
1✔
361
        (
362
            "Fetch data async (measure: %s, host: %s, field: %s, "
363
            "start: %s, stop: %s, every: %dmin, window: %dmin, "
364
            "create_empty: %s, last: %s)"
365
        ),
366
        measure,
367
        hostname,
368
        field,
369
        start,
370
        stop,
371
        every_min,
372
        window_min,
373
        create_empty,
374
        last,
375
    )
376

377
    try:
1✔
378
        template = FLUX_QUERY_WITHOUT_AGGREGATION if window_min == 0 else FLUX_QUERY
1✔
379

380
        table_list = await _fetch_data_impl_async(
1✔
381
            db_config,
382
            template,
383
            measure,
384
            hostname,
385
            field,
386
            start,
387
            stop,
388
            every_min,
389
            window_min,
390
            create_empty,
391
            last,
392
        )
393
        time_fetched = time.time()
1✔
394

395
        result = _process_query_results(table_list, create_empty, last, every_min, window_min)
1✔
396

397
        time_finish = time.time()
1✔
398
        if ((time_fetched - time_start) > 1) or ((time_finish - time_fetched) > 0.1):
1✔
399
            logging.warning(
×
400
                "It's taking too long to retrieve the data. (fetch: %.2f sec, modify: %.2f sec)",
401
                time_fetched - time_start,
402
                time_finish - time_fetched,
403
            )
404

405
        return result
1✔
NEW
406
    except Exception as e:
×
407
        logging.exception("Failed to fetch data")
×
408

NEW
409
        return SensorDataResult(error_message=str(e))
×
410

411

412
async def fetch_data_parallel(
1✔
413
    db_config: InfluxDBConfig, requests: list[DataRequest]
414
) -> list[SensorDataResult | BaseException]:
415
    """
416
    複数のデータ取得リクエストを並列実行
417

418
    Args:
419
    ----
420
        db_config: InfluxDBの設定(全リクエスト共通)
421
        requests: DataRequest のリスト
422

423
    Returns:
424
    -------
425
        各リクエストの結果を含むリスト
426

427
    """
428
    tasks = []
1✔
429
    for req in requests:
1✔
430
        task = fetch_data_async(
1✔
431
            db_config,
432
            req.measure,
433
            req.hostname,
434
            req.field,
435
            req.start,
436
            req.stop,
437
            req.every_min,
438
            req.window_min,
439
            req.create_empty,
440
            req.last,
441
        )
442
        tasks.append(task)
1✔
443

444
    return await asyncio.gather(*tasks, return_exceptions=True)
1✔
445

446

447
def get_equip_on_minutes(  # noqa: PLR0913
1✔
448
    config: InfluxDBConfig,
449
    measure: str,
450
    hostname: str,
451
    field: str,
452
    threshold: float,
453
    start: str = "-30h",
454
    stop: str = "now()",
455
    every_min: int = 1,
456
    window_min: int = 5,
457
    create_empty: bool = True,
458
) -> int:
459
    logging.info(
1✔
460
        (
461
            "Get 'ON' minutes (type: %s, host: %d, field: %d{field}, "
462
            "threshold: %.2f, start: %s, stop: %s, every: %smin, "
463
            "window: %dmin, create_empty: %s)"
464
        ),
465
        measure,
466
        hostname,
467
        field,
468
        threshold,
469
        start,
470
        stop,
471
        every_min,
472
        window_min,
473
        create_empty,
474
    )
475

476
    try:
1✔
477
        table_list = _fetch_data_impl(
1✔
478
            config,
479
            FLUX_QUERY,
480
            measure,
481
            hostname,
482
            field,
483
            start,
484
            stop,
485
            every_min,
486
            window_min,
487
            create_empty,
488
        )
489

490
        if len(table_list) == 0:
1✔
491
            return 0
1✔
492

493
        count = 0
1✔
494

495
        every_min = int(every_min)
1✔
496
        window_min = int(window_min)
1✔
497
        record_num = len(table_list[0].records)
1✔
498
        for i, record in enumerate(table_list[0].records):
1✔
499
            if create_empty and (window_min > every_min) and (i > record_num - 1 - (window_min - every_min)):
1✔
500
                # NOTE: timedMovingAverage を使うと、末尾に余分なデータが入るので取り除く
501
                continue
×
502

503
            # NOTE: aggregateWindow(createEmpty: true) と fill(usePrevious: true) の組み合わせ
504
            # だとタイミングによって、先頭に None が入る
505
            if record.get_value() is None:
1✔
506
                continue
×
507
            if record.get_value() >= threshold:
1✔
508
                count += 1
1✔
509

510
        return count * int(every_min)
1✔
511
    except Exception:
1✔
512
        logging.exception("Failed to fetch data")
1✔
513
        return 0
1✔
514

515

516
def get_equip_mode_period(  # noqa: C901, PLR0913
1✔
517
    config: InfluxDBConfig,
518
    measure: str,
519
    hostname: str,
520
    field: str,
521
    threshold_list: list[float],
522
    start: str = "-30h",
523
    stop: str = "now()",
524
    every_min: int = 10,
525
    window_min: int = 10,
526
    create_empty: bool = True,
527
) -> list[list[Any]]:
528
    logging.info(
×
529
        (
530
            "Get equipment mode period (type: %s, host: %s, field: %s, "
531
            "threshold: %.2f, start: %s, stop: %s, every: %dmin, "
532
            "window: %dmin, create_empty: %s)",
533
        ),
534
        measure,
535
        hostname,
536
        field,
537
        f"[{','.join(f'{v:.1f}' for v in threshold_list)}]",
538
        start,
539
        stop,
540
        every_min,
541
        window_min,
542
        create_empty,
543
    )
544

545
    try:
×
546
        table_list = _fetch_data_impl(
×
547
            config,
548
            FLUX_QUERY,
549
            measure,
550
            hostname,
551
            field,
552
            start,
553
            stop,
554
            every_min,
555
            window_min,
556
            create_empty,
557
        )
558

559
        if len(table_list) == 0:
×
560
            return []
×
561

562
        # NOTE: 常時冷却と間欠冷却の期間を求める
563
        on_range = []
×
564
        state = -1
×
565
        start_time = None
×
566
        prev_time = None
×
567
        localtime_offset = datetime.timedelta(hours=9)
×
568

569
        for record in table_list[0].records:
×
570
            # NOTE: aggregateWindow(createEmpty: true) と fill(usePrevious: true) の組み合わせ
571
            # だとタイミングによって、先頭に None が入る
572
            if record.get_value() is None:
×
573
                logging.debug("DELETE %s", record.get_time() + localtime_offset)
×
574
                continue
×
575

576
            is_idle = True
×
577
            for i in range(len(threshold_list)):
×
578
                if record.get_value() > threshold_list[i]:
×
579
                    if state != i:
×
580
                        if state != -1:
×
581
                            assert start_time is not None
×
582
                            assert prev_time is not None
×
583
                            on_range.append(
×
584
                                [
585
                                    start_time + localtime_offset,
586
                                    prev_time + localtime_offset,
587
                                    state,
588
                                ]
589
                            )
590
                        state = i
×
591
                        start_time = record.get_time()
×
592
                    is_idle = False
×
593
                    break
×
594
            if is_idle and state != -1:
×
595
                assert start_time is not None
×
596
                assert prev_time is not None
×
597
                on_range.append(
×
598
                    [
599
                        start_time + localtime_offset,
600
                        prev_time + localtime_offset,
601
                        state,
602
                    ]
603
                )
604
                state = -1
×
605
                start_time = record.get_time()
×
606

607
            prev_time = record.get_time()
×
608

609
        if state != -1:
×
610
            assert start_time is not None
×
611
            on_range.append(
×
612
                [
613
                    start_time + localtime_offset,
614
                    table_list[0].records[-1].get_time() + localtime_offset,
615
                    state,
616
                ]
617
            )
618
        return on_range
×
619
    except Exception:
×
620
        logging.exception("Failed to fetch data")
×
621
        return []
×
622

623

624
def get_sum(  # noqa: PLR0913
1✔
625
    config: InfluxDBConfig,
626
    measure: str,
627
    hostname: str,
628
    field: str,
629
    start: str = "-3m",
630
    stop: str = "now()",
631
    every_min: int = 1,
632
    window_min: int = 3,
633
) -> float:
634
    try:
1✔
635
        table_list = _fetch_data_impl(
1✔
636
            config, FLUX_SUM_QUERY, measure, hostname, field, start, stop, every_min, window_min, True
637
        )
638

639
        value_list = table_list.to_values(columns=["count", "sum"])
1✔
640

641
        if len(value_list) == 0:
1✔
642
            return 0
1✔
643
        else:
644
            sum_value = value_list[0][1]
1✔
645
            if isinstance(sum_value, (int, float)):
1✔
646
                return float(sum_value)
1✔
647
            logging.warning("Unexpected sum value type: %s (value=%s)", type(sum_value).__name__, sum_value)
×
648
            return 0
×
649
    except Exception:
1✔
650
        logging.exception("Failed to fetch data")
1✔
651
        return 0
1✔
652

653

654
def get_day_sum(  # noqa: PLR0913
1✔
655
    config: InfluxDBConfig,
656
    measure: str,
657
    hostname: str,
658
    field: str,
659
    days: int,
660
    day_before: int = 0,
661
    day_offset: int = 0,
662
    every_min: int = 1,
663
    window_min: int = 5,
664
) -> float:
665
    now = my_lib.time.now()
1✔
666

667
    if day_before == 0:
1✔
668
        start = f"-{day_offset + days - 1}d{now.hour}h{now.minute}m"
1✔
669
        stop = f"-{day_offset}d"
1✔
670
    else:
671
        start = f"-{day_before + day_offset + days - 1}d{now.hour}h{now.minute}m"
×
672
        stop = f"-{day_before + day_offset - 1}d{now.hour}h{now.minute}m"
×
673

674
    return get_sum(config, measure, hostname, field, start, stop, every_min, window_min)
1✔
675

676

677
def get_hour_sum(  # noqa: PLR0913
1✔
678
    config: InfluxDBConfig,
679
    measure: str,
680
    hostname: str,
681
    field: str,
682
    hours: int,
683
    day_offset: int = 0,
684
    every_min: int = 1,
685
    window_min: int = 1,
686
) -> float:
687
    start = f"-{day_offset * 24 + hours}h"
1✔
688
    stop = f"-{day_offset * 24}h"
1✔
689

690
    return get_sum(config, measure, hostname, field, start, stop, every_min, window_min)
1✔
691

692

693
def get_minute_sum(  # noqa: PLR0913
1✔
694
    config: InfluxDBConfig,
695
    measure: str,
696
    hostname: str,
697
    field: str,
698
    minutes: int,
699
    day_offset: int = 0,
700
    every_min: int = 1,
701
    window_min: int = 1,
702
) -> float:
703
    start = f"-{day_offset * 24 * 60 + minutes}m"
1✔
704
    stop = f"-{day_offset * 24 * 60}m"
1✔
705

706
    return get_sum(config, measure, hostname, field, start, stop, every_min, window_min)
1✔
707

708

709
def get_last_event(
1✔
710
    config: InfluxDBConfig, measure: str, hostname: str, field: str, start: str = "-7d"
711
) -> datetime.datetime | None:
712
    try:
1✔
713
        table_list = _fetch_data_impl(
1✔
714
            config, FLUX_EVENT_QUERY, measure, hostname, field, start, "now()", 0, 0, False
715
        )
716

717
        value_list = table_list.to_values(columns=["_time"])
1✔
718

719
        if len(value_list) == 0:
1✔
720
            return None
1✔
721
        else:
722
            time_value = value_list[0][0]
1✔
723
            if isinstance(time_value, datetime.datetime):
1✔
724
                return time_value
1✔
725
            logging.warning(
×
726
                "Unexpected time value type: %s (value=%s)", type(time_value).__name__, time_value
727
            )
728
            return None
×
729
    except Exception:
1✔
730
        logging.exception("Failed to fetch data")
1✔
731
        return None
1✔
732

733

734
def dump_data(data: SensorDataResult) -> None:
1✔
735
    for i in range(len(data.time)):
1✔
736
        logging.info("%s: %s", data.time[i], data.value[i])
1✔
737

738

739
if __name__ == "__main__":
740
    # TEST Code
741
    import docopt
742

743
    import my_lib.config
744
    import my_lib.logger
745
    import my_lib.pretty
746

747
    def get_config(config, dotted_key):
748
        keys = dotted_key.split(".")
749
        value = config
750

751
        for key in keys:
752
            value = value[key]
753

754
        return value
755

756
    assert __doc__ is not None
757
    args = docopt.docopt(__doc__)
758

759
    config_file = args["-c"]
760
    mode = args["-m"]
761
    every = args["-e"]
762
    window = args["-w"]
763
    infxlux_db_spec = args["-i"]
764
    sensor_spec = args["-s"]
765
    field_name = args["-f"]
766
    period = int(args["-p"])
767
    debug_mode = args["-D"]
768

769
    my_lib.logger.init("test", level=logging.DEBUG if debug_mode else logging.INFO)
770

771
    config = my_lib.config.load(config_file)
772

773
    db_config = cast(InfluxDBConfig, get_config(config, infxlux_db_spec))
774
    sensor_config = get_config(config, sensor_spec)
775

776
    logging.info("DB config: %s", my_lib.pretty.format(db_config))
777
    logging.info("Sensor config: %s", my_lib.pretty.format(sensor_config))
778

779
    result: SensorDataResult | float
780
    if mode == "data":
781
        result = fetch_data(
782
            db_config,
783
            sensor_config["measure"],
784
            sensor_config["hostname"],
785
            field_name,
786
            start="-10m",
787
            stop="now()",
788
            every_min=1,
789
            window_min=3,
790
            create_empty=True,
791
            last=False,
792
        )
793
    elif mode == "day_sum":
794
        result = get_day_sum(
795
            db_config, sensor_config["measure"], sensor_config["hostname"], field_name, period
796
        )
797
    elif mode == "hour_sum":
798
        result = get_hour_sum(
799
            db_config, sensor_config["measure"], sensor_config["hostname"], field_name, period
800
        )
801
    elif mode == "minute_sum":
802
        result = get_minute_sum(
803
            db_config, sensor_config["measure"], sensor_config["hostname"], field_name, period
804
        )
805
    else:
806
        logging.error("Unknown mode: %s", mode)
807
        result = 0.0
808

809
    logging.info(my_lib.pretty.format(result))
STATUS · Troubleshooting · Open an Issue · Sales · Support · CAREERS · ENTERPRISE · START FREE TRIAL · SCHEDULE DEMO
ANNOUNCEMENTS · TWITTER · TOS & SLA · Supported CI Services · What's a CI service? · Automated Testing

© 2026 Coveralls, Inc