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

kimata / my-py-lib / 20691871150

04 Jan 2026 10:58AM UTC coverage: 64.546% (-0.2%) from 64.783%
20691871150

push

github

kimata
feat: コンテナの起動経過時間を取得する container_util.get_uptime() を追加

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

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

115 of 157 new or added lines in 37 files covered. (73.25%)

4 existing lines in 3 files now uncovered.

3095 of 4795 relevant lines covered (64.55%)

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

22
from __future__ import annotations
1✔
23

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

32

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

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

41

42
import influxdb_client  # noqa: E402
1✔
43
from influxdb_client.client.flux_table import TableList  # noqa: E402
1✔
44

45
import my_lib.time  # noqa: E402
1✔
46

47

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

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

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

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

83

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

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

98

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

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

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

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

150

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

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

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

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

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

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

203

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

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

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

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

249

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

281

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

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

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

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

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

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

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

346

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

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

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

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

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

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

410
        return SensorDataResult(error_message=str(e))
×
411

412

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

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

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

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

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

447

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

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

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

494
        count = 0
1✔
495

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

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

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

516

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

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

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

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

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

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

608
            prev_time = record.get_time()
×
609

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

624

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

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

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

654

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

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

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

677

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

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

693

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

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

709

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

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

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

734

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

739

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

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

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

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

755
        return value
756

757
    assert __doc__ is not None  # noqa: S101
758
    args = docopt.docopt(__doc__)
759

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

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

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

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

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

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

810
    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