Fix ansible-lint violations: FQCN, formatting, bugs, role renames
- Auto-fix FQCN, YAML formatting, jinja spacing, and free-form module syntax via ansible-lint --fix - Fix comments misplaced inside module args by the auto-fixer (bluetooth-monitor, pi_standard_setup, pi_musicmouse) - Fix notify: references left stale (lowercase) after handler names were re-cased, which would have silently broken reboot/restart handlers (pi_disable_onboard_bluetooth, pi_hifiberry_amp, pi_squeezelite, pi_standard_setup) - Fix a task in pis/debmatic-install.yml missing its module name (apt_repository), which caused a real syntax-check failure - Add missing play names, fix comment spacing, literal-compare idiom, and no-changed-when annotations - Delete unused/broken roles/better-shell-env (unreferenced, invalid YAML) - Rename all hyphenated role directories to underscore form to satisfy ansible-lint's role-name rule, updating every playbook/meta reference Remaining lint findings (var-naming, package-latest, risky-file-permissions, no-handler) intentionally left for follow-up per user decision.
This commit is contained in:
12
roles/bluetooth_monitor/README.md
Normal file
12
roles/bluetooth_monitor/README.md
Normal file
@@ -0,0 +1,12 @@
|
||||
# bluetooth_monitor
|
||||
|
||||
Installs `my_btmonitor.py` (BLE scanning via `bleak`) as a systemd service
|
||||
that watches for nearby devices, publishes state over MQTT, and can restart
|
||||
the local BLE interface on a watchdog timeout.
|
||||
|
||||
**Key vars:** `my_bt_monitor_watchdog_seconds`, `my_btmonitor_restart_ble_interface`,
|
||||
`my_btmonitor_mqtt_username`, `my_btmonitor_mqtt_password`
|
||||
|
||||
`other/` is not part of the deployed role — it's a separate, standalone
|
||||
data-analysis project (Jupyter notebook, Dockerfile, collected CSV data) used
|
||||
to analyze data captured by the monitor.
|
||||
131
roles/bluetooth_monitor/files/filter.cpp
Normal file
131
roles/bluetooth_monitor/files/filter.cpp
Normal file
@@ -0,0 +1,131 @@
|
||||
#include <cmath>
|
||||
#include <vector>
|
||||
#include <iostream>
|
||||
|
||||
using real_t = double;
|
||||
|
||||
static constexpr real_t SPIKE_THRESHOLD = 1.0f; // Threshold for spike detection
|
||||
static constexpr int NUM_READINGS = 12; // Number of readings to keep track of
|
||||
|
||||
|
||||
class FilteredDistance {
|
||||
public:
|
||||
FilteredDistance(real_t minCutoff = 1e-1f, real_t beta = 1e-3, real_t dcutoff = 5e-3);
|
||||
void addMeasurement(real_t dist, real_t time_now_in_seconds);
|
||||
const real_t getMedianDistance() const;
|
||||
const real_t getDistance() const;
|
||||
const real_t getVariance() const;
|
||||
|
||||
bool hasValue() const { return lastTime != 0; }
|
||||
|
||||
private:
|
||||
real_t minCutoff;
|
||||
real_t beta;
|
||||
real_t dcutoff;
|
||||
real_t x, dx;
|
||||
real_t lastDist;
|
||||
real_t lastTime;
|
||||
|
||||
real_t getAlpha(real_t cutoff, real_t dT);
|
||||
|
||||
real_t readings[NUM_READINGS]; // Array to store readings
|
||||
int readIndex; // Current position in the array
|
||||
real_t total; // Total of the readings
|
||||
real_t totalSquared; // Total of the squared readings
|
||||
|
||||
void initSpike(real_t dist);
|
||||
real_t removeSpike(real_t dist);
|
||||
};
|
||||
|
||||
FilteredDistance::FilteredDistance(real_t minCutoff, real_t beta, real_t dcutoff)
|
||||
: minCutoff(minCutoff), beta(beta), dcutoff(dcutoff), x(0), dx(0), lastDist(0), lastTime(-1), total(0), totalSquared(0), readIndex(0) {
|
||||
}
|
||||
|
||||
void FilteredDistance::initSpike(real_t dist) {
|
||||
for (size_t i = 0; i < NUM_READINGS; i++) {
|
||||
readings[i] = dist;
|
||||
}
|
||||
total = dist * NUM_READINGS;
|
||||
totalSquared = dist * dist * NUM_READINGS; // Initialize sum of squared distances
|
||||
}
|
||||
|
||||
real_t FilteredDistance::removeSpike(real_t dist) {
|
||||
total -= readings[readIndex]; // Subtract the last reading
|
||||
totalSquared -= readings[readIndex] * readings[readIndex]; // Subtract the square of the last reading
|
||||
|
||||
readings[readIndex] = dist; // Read the sensor
|
||||
total += readings[readIndex]; // Add the reading to the total
|
||||
totalSquared += readings[readIndex] * readings[readIndex]; // Add the square of the reading
|
||||
|
||||
readIndex = (readIndex + 1) % NUM_READINGS; // Advance to the next position in the array
|
||||
|
||||
auto average = total / static_cast<real_t>(NUM_READINGS); // Calculate the average
|
||||
|
||||
if (std::fabs(dist - average) > SPIKE_THRESHOLD)
|
||||
return average; // Spike detected, use the average as the filtered value
|
||||
|
||||
return dist; // No spike, return the new value
|
||||
}
|
||||
|
||||
void FilteredDistance::addMeasurement(real_t dist, real_t time_now_in_seconds) {
|
||||
const bool initialized = lastTime >= 0;
|
||||
const real_t elapsed = time_now_in_seconds - lastTime;
|
||||
lastTime = time_now_in_seconds;
|
||||
|
||||
if (!initialized) {
|
||||
x = dist; // Set initial filter state to the first reading
|
||||
dx = 0; // Initial derivative is unknown, so we set it to zero
|
||||
lastDist = dist;
|
||||
initSpike(dist);
|
||||
} else {
|
||||
real_t dT = std::max(elapsed, real_t(0.05)); // Convert microseconds to seconds, enforce a minimum dT
|
||||
const real_t alpha = getAlpha(minCutoff, dT);
|
||||
const real_t dAlpha = getAlpha(dcutoff, dT);
|
||||
dist = removeSpike(dist);
|
||||
x += alpha * (dist - x);
|
||||
dx = dAlpha * ((dist - lastDist) / dT);
|
||||
lastDist = x + beta * dx;
|
||||
std::cout << "alpha=" << alpha <<
|
||||
" dAlpha=" << dAlpha <<
|
||||
" dist=" << dist <<
|
||||
" x=" << x <<
|
||||
" dx=" << dx <<
|
||||
" lastDist=" << lastDist <<
|
||||
std::endl;
|
||||
}
|
||||
}
|
||||
|
||||
const real_t FilteredDistance::getDistance() const {
|
||||
return lastDist;
|
||||
}
|
||||
|
||||
real_t FilteredDistance::getAlpha(real_t cutoff, real_t dT) {
|
||||
real_t tau = 1.0f / (2 * M_PI * cutoff);
|
||||
return 1.0f / (1.0f + tau / dT);
|
||||
}
|
||||
|
||||
const real_t FilteredDistance::getVariance() const {
|
||||
auto mean = total / static_cast<real_t>(NUM_READINGS);
|
||||
auto meanOfSquares = totalSquared / static_cast<real_t>(NUM_READINGS);
|
||||
auto variance = meanOfSquares - (mean * mean); // Variance formula: E(X^2) - (E(X))^2
|
||||
if (variance < 0.0f) return 0.0f;
|
||||
return variance;
|
||||
}
|
||||
|
||||
|
||||
int main(int argc, char**argv)
|
||||
{
|
||||
FilteredDistance f;
|
||||
std::vector<real_t> values = {1.5, 2.9, 5.3, 15.1, 1.5, 2.5, 1.5, 2.9, 5.3, 15.1};
|
||||
|
||||
real_t time = 0.0;
|
||||
//std::cout << " result_cpp = [";
|
||||
for(int i=0; i < 1; ++i)
|
||||
for(auto value : values) {
|
||||
f.addMeasurement(value, time);
|
||||
time += 1.0;
|
||||
//std::cout << f.getDistance() << ", ";
|
||||
}
|
||||
//std::cout << "]" << std::endl;
|
||||
return 0;
|
||||
}
|
||||
91
roles/bluetooth_monitor/files/filter.py
Normal file
91
roles/bluetooth_monitor/files/filter.py
Normal file
@@ -0,0 +1,91 @@
|
||||
#from time import time
|
||||
import math
|
||||
from scipy import signal
|
||||
|
||||
# Taken from ESPresense C++ code
|
||||
class FilteredDistance:
|
||||
NUM_READINGS = 100
|
||||
SPIKE_THRESHOLD = 1.0
|
||||
|
||||
def __init__(self, min_cutoff : float = 1e-1, beta : float = 1e-3, dcutoff : float = 5e-3):
|
||||
self.min_cutoff = min_cutoff
|
||||
self.beta = beta
|
||||
self.dcutoff = dcutoff
|
||||
self.x = 0
|
||||
self.dx = 0
|
||||
self.last_dist = 0
|
||||
self.last_time = -1.0
|
||||
self.total = 0
|
||||
self.read_index = 0
|
||||
self.readings = []
|
||||
|
||||
def _init_spike(self, dist : float):
|
||||
self.readings = [dist] * self.NUM_READINGS
|
||||
self.total = sum(self.readings)
|
||||
|
||||
def _remove_spike(self, dist: float):
|
||||
self.total -= self.readings[self.read_index]
|
||||
|
||||
self.readings[self.read_index] = dist
|
||||
|
||||
self.total += dist
|
||||
|
||||
self.read_index = (self.read_index + 1) % self.NUM_READINGS
|
||||
average = self.total / self.NUM_READINGS
|
||||
if abs(dist - average) > self.SPIKE_THRESHOLD:
|
||||
return average # spike detected
|
||||
else:
|
||||
return dist
|
||||
|
||||
def _get_alpha(self, cutoff : float, dT : float):
|
||||
tau = 1 / (2 * math.pi * cutoff)
|
||||
return 1 / (1 + tau / dT)
|
||||
|
||||
def add_measurement(self, dist : float, time_now_in_seconds: float):
|
||||
initialized = (self.last_time >= 0.0)
|
||||
elapsed = time_now_in_seconds - self.last_time
|
||||
self.last_time = time_now_in_seconds
|
||||
if not initialized:
|
||||
self.x = dist
|
||||
self.dx = 0
|
||||
self.last_dist = dist
|
||||
self._init_spike(dist)
|
||||
else:
|
||||
dT = max(elapsed, 0.05)
|
||||
alpha = self._get_alpha(self.min_cutoff, dT)
|
||||
d_alpha = self._get_alpha(self.dcutoff, dT)
|
||||
dist = self._remove_spike(dist)
|
||||
self.x += alpha * (dist - self.x)
|
||||
self.dx = d_alpha * ((dist - self.last_dist) / dT)
|
||||
self.last_dist = self.x + self.beta * self.dx
|
||||
#print(f"{alpha=} {d_alpha=} {dist=} {self.x=} {self.dx=} {self.last_dist=}")
|
||||
|
||||
def get_distance(self):
|
||||
return self.last_dist
|
||||
|
||||
def run_test(times, values, **kwargs):
|
||||
f = FilteredDistance(**kwargs)
|
||||
result = []
|
||||
for t, value in zip(times, values):
|
||||
f.add_measurement(value, t)
|
||||
result.append(f.get_distance())
|
||||
return result
|
||||
|
||||
def smooth(y, box_pts):
|
||||
box = np.ones(box_pts)/box_pts
|
||||
y_smooth = np.convolve(y, box, mode='same')
|
||||
return y_smooth
|
||||
|
||||
if __name__ == "__main__":
|
||||
import numpy as np
|
||||
values = np.array([1] * 20 + [2, 4, 6, 7, 10, 16, 10, 13, 16, 24, 13] + [1] * 20 )
|
||||
|
||||
times = np.arange(0, len(values)) * 10
|
||||
result_default = run_test(times, values)
|
||||
result_beta1 = smooth(values, 6)
|
||||
import matplotlib.pyplot as plt
|
||||
plt.plot(times, values, label="raw")
|
||||
#plt.plot(times, result_default, marker="o", label="filtered")
|
||||
plt.plot(times, result_beta1, marker='x', label="altered")
|
||||
plt.legend()
|
||||
plt.show()
|
||||
BIN
roles/bluetooth_monitor/files/filtered
Executable file
BIN
roles/bluetooth_monitor/files/filtered
Executable file
Binary file not shown.
11
roles/bluetooth_monitor/files/my_btmonitor.service
Normal file
11
roles/bluetooth_monitor/files/my_btmonitor.service
Normal file
@@ -0,0 +1,11 @@
|
||||
[Unit]
|
||||
Description=My Bluetooth monitor
|
||||
After=network.target
|
||||
|
||||
[Service]
|
||||
Type=simple
|
||||
Restart=always
|
||||
ExecStart=/usr/bin/my_btmonitor
|
||||
|
||||
[Install]
|
||||
WantedBy=multi-user.target
|
||||
11
roles/bluetooth_monitor/other/Dockerfile
Normal file
11
roles/bluetooth_monitor/other/Dockerfile
Normal file
@@ -0,0 +1,11 @@
|
||||
FROM python:3
|
||||
|
||||
WORKDIR /usr/src/app
|
||||
|
||||
COPY requirements.txt ./
|
||||
RUN pip install --no-cache-dir -r requirements.txt
|
||||
|
||||
COPY bt_monitor_server.py .
|
||||
COPY training_data.csv .
|
||||
|
||||
CMD [ "python", "./bt_monitor_server.py" ]
|
||||
95
roles/bluetooth_monitor/other/analysis.py
Normal file
95
roles/bluetooth_monitor/other/analysis.py
Normal file
@@ -0,0 +1,95 @@
|
||||
from pathlib import Path
|
||||
import pandas as pd
|
||||
import numpy as np
|
||||
from copy import copy
|
||||
from sklearn.model_selection import cross_val_score
|
||||
from sklearn import svm
|
||||
from sklearn.neural_network import MLPClassifier
|
||||
from sklearn.metrics import confusion_matrix, ConfusionMatrixDisplay
|
||||
import matplotlib.pyplot as plt
|
||||
from sklearn.model_selection import train_test_split
|
||||
|
||||
def load_measurements(csv_file: Path):
|
||||
def cleanup_column_name(col_name: str):
|
||||
clean_name = col_name.replace('#', '').strip()
|
||||
if clean_name == 'room':
|
||||
return 'tracker'
|
||||
return clean_name
|
||||
|
||||
df = pd.read_csv(str(csv_file))
|
||||
|
||||
# String cleanup in column names and room names
|
||||
df = df.rename(columns=cleanup_column_name)
|
||||
df.applymap(lambda x: x.strip() if isinstance(x, str) else x)
|
||||
|
||||
df['tracker'] = df['tracker'].astype("category")
|
||||
df['real_room'] = df['real_room'].astype("category")
|
||||
|
||||
return df
|
||||
|
||||
|
||||
FAR_AWAY_FEATURE_VALUE = 1
|
||||
def get_feature_value(rssi, tx_power):
|
||||
MIN_RSSI = -90
|
||||
MAX_TRANSFORMED_RSSI = 40
|
||||
v = tx_power - rssi - MAX_TRANSFORMED_RSSI
|
||||
if v < 0:
|
||||
v = 0
|
||||
return v / (-MIN_RSSI)
|
||||
|
||||
|
||||
def make_training_data(df: pd.DataFrame, device_to_map):
|
||||
idx_to_tracker = dict(enumerate(df['tracker'].cat.categories ))
|
||||
tracker_to_idx = {v: k for k, v in idx_to_tracker.items()}
|
||||
idx_to_room = dict(enumerate(df['real_room'].cat.categories ))
|
||||
room_to_idx = {v: k for k, v in idx_to_room.items()}
|
||||
|
||||
last_real_room = None
|
||||
start_time = None
|
||||
current_feature = [FAR_AWAY_FEATURE_VALUE] * len(idx_to_tracker)
|
||||
|
||||
features = []
|
||||
labels = []
|
||||
|
||||
# Feature vectors - rssi column for each room
|
||||
for i, row in df.iterrows():
|
||||
time, device, tracker, rssi, tx_power, real_room = row
|
||||
if device != device_to_map:
|
||||
continue
|
||||
if last_real_room != real_room:
|
||||
start_time = time
|
||||
last_real_room = real_room
|
||||
|
||||
tracker_idx = tracker_to_idx[tracker]
|
||||
current_feature[tracker_idx] = get_feature_value(rssi, tx_power)
|
||||
if time - start_time > 20:
|
||||
features.append(copy(current_feature))
|
||||
labels.append(room_to_idx[real_room])
|
||||
|
||||
return np.array(features), np.array(labels)
|
||||
|
||||
def train(features, labels, classes):
|
||||
clf = svm.SVC(kernel='rbf')
|
||||
print("Training")
|
||||
scores = cross_val_score(clf, features, labels, cv=5)
|
||||
print(scores)
|
||||
print("%0.2f accuracy with a standard deviation of %0.2f" % (scores.mean(), scores.std()))
|
||||
|
||||
X_train, X_test, y_train, y_test = train_test_split(features, labels, random_state=0)
|
||||
clf.fit(X_train, y_train)
|
||||
cm = confusion_matrix(clf.predict(X_test), y_test)
|
||||
print(cm)
|
||||
print(classes)
|
||||
disp = ConfusionMatrixDisplay(confusion_matrix=cm, display_labels=classes)
|
||||
disp.plot()
|
||||
plt.show()
|
||||
|
||||
|
||||
|
||||
if __name__ == "__main__":
|
||||
csv_path = Path("/home/martin/code/ansible/roles/bluetooth-monitor/other/collected.csv")
|
||||
df = load_measurements(csv_path)
|
||||
features, labels = make_training_data(df, "martins_apple_watch")
|
||||
print(np.unique(labels))
|
||||
print(features.shape, labels.shape)
|
||||
train(features, labels, list(df['real_room'].dtype.categories))
|
||||
383
roles/bluetooth_monitor/other/bt_monitor_analyze.ipynb
Normal file
383
roles/bluetooth_monitor/other/bt_monitor_analyze.ipynb
Normal file
File diff suppressed because one or more lines are too long
309
roles/bluetooth_monitor/other/bt_monitor_server.py
Executable file
309
roles/bluetooth_monitor/other/bt_monitor_server.py
Executable file
@@ -0,0 +1,309 @@
|
||||
#!/usr/bin/env python3
|
||||
|
||||
import os
|
||||
import aiomqtt
|
||||
import json
|
||||
import asyncio
|
||||
from time import time
|
||||
from pathlib import Path
|
||||
from collections import namedtuple, defaultdict, deque
|
||||
from typing import Dict, Optional, List
|
||||
from Crypto.Cipher import AES
|
||||
import pandas as pd
|
||||
import numpy as np
|
||||
import logging
|
||||
from sklearn import svm
|
||||
from sklearn.model_selection import cross_val_score
|
||||
|
||||
logging.basicConfig(level=logging.INFO)
|
||||
|
||||
BtleMeasurement = namedtuple("BtleMeasurement", ["time", "tracker", "address", "rssi", "tx_power"])
|
||||
BtleDeviceMeasurement = namedtuple("BtleDeviceMeasurement", ["time", "device", "tracker", "rssi", "tx_power"])
|
||||
MqttInfo = namedtuple("MqttInfo", ["server", "username", "password"])
|
||||
|
||||
# ------------------------------------------------------- DECODING -------------------------------------------------------------------------
|
||||
|
||||
|
||||
class DeviceDecoder:
|
||||
"""Decode bluetooth addresses - either simple ones (just address to name) or random changing ones like Apple devices using irk keys"""
|
||||
|
||||
def __init__(self, irk_to_devicename: Dict[str, str], address_to_name: Dict[str, str]):
|
||||
"""
|
||||
address_to_name: dictionary from bt address as string separated by ":" to a device name
|
||||
irk_to_devicename is dict with irk as a hex string, mapping to device name
|
||||
"""
|
||||
self.irk_to_devicename = {bytes.fromhex(k): v for k, v in irk_to_devicename.items()}
|
||||
self.address_to_name = address_to_name
|
||||
|
||||
def _resolve_rpa(rpa: bytes, irk: bytes) -> bool:
|
||||
"""Compares the random address rpa to an irk (secret key) and return True if it matches"""
|
||||
assert len(rpa) == 6
|
||||
assert len(irk) == 16
|
||||
|
||||
key = irk
|
||||
plain_text = b"\x00" * 16
|
||||
plain_text = bytearray(plain_text)
|
||||
plain_text[15] = rpa[3]
|
||||
plain_text[14] = rpa[4]
|
||||
plain_text[13] = rpa[5]
|
||||
plain_text = bytes(plain_text)
|
||||
|
||||
cipher = AES.new(key, AES.MODE_ECB)
|
||||
cipher_text = cipher.encrypt(plain_text)
|
||||
return cipher_text[15] == rpa[0] and cipher_text[14] == rpa[1] and cipher_text[13] == rpa[2]
|
||||
|
||||
def _addr_to_bytes(addr: str) -> bytes:
|
||||
"""Converts a bluetooth mac address string with semicolons to bytes"""
|
||||
str_without_colons = addr.replace(":", "")
|
||||
bytearr = bytearray.fromhex(str_without_colons)
|
||||
bytearr.reverse()
|
||||
return bytes(bytearr)
|
||||
|
||||
def decode(self, addr: str) -> Optional[str]:
|
||||
"""addr is a bluetooth address as a string e.g. 4d:24:12:12:34:10"""
|
||||
for irk, name in self.irk_to_devicename.items():
|
||||
if DeviceDecoder._resolve_rpa(DeviceDecoder._addr_to_bytes(addr), irk):
|
||||
return name
|
||||
return self.address_to_name.get(addr, None)
|
||||
|
||||
def __call__(self, m: BtleMeasurement) -> Optional[BtleDeviceMeasurement]:
|
||||
decoded_device_name = self.decode(m.address)
|
||||
if not decoded_device_name:
|
||||
return None
|
||||
return BtleDeviceMeasurement(m.time, decoded_device_name, m.tracker, m.rssi, m.tx_power)
|
||||
|
||||
|
||||
# ------------------------------------------------------- MACHINE LEARNING ----------------------------------------------------------------
|
||||
|
||||
|
||||
class KnownRoomCsvLogger:
|
||||
"""Logs known room measurements to be used later as training data for classifier"""
|
||||
|
||||
def __init__(self, csv_file: Path):
|
||||
self.known_room = None
|
||||
|
||||
if csv_file.exists():
|
||||
self.csv_file_handle = open(csv_file, "a")
|
||||
else:
|
||||
self.csv_file_handle = open(csv_file, "w")
|
||||
print(f"#time,device,tracker,rssi,tx_power,known_room", file=csv_file)
|
||||
|
||||
def update_known_room(self, known_room: str):
|
||||
if known_room != self.known_room:
|
||||
logging.info(f"Updating known_room {self.known_room} -> {known_room}")
|
||||
self.known_room = known_room
|
||||
|
||||
def report_measure(self, m: BtleDeviceMeasurement):
|
||||
ignore_rooms = ("keins", "?", "none", "unknown")
|
||||
if self.known_room is None or self.known_room in ignore_rooms:
|
||||
return
|
||||
logging.info(f"Appending to training set: {m}")
|
||||
print(
|
||||
f"{m.time},{m.device},{m.tracker},{m.rssi},{m.tx_power},{self.known_room}",
|
||||
file=self.csv_file_handle,)
|
||||
|
||||
|
||||
class RunningFeatureVector:
|
||||
FAR_AWAY_FEATURE_VALUE = 1
|
||||
MIN_TIME_UNTIL_PREDICTION = 40 # wait until every reachable tracker detected the device
|
||||
TIME_TO_DELETE_IF_NOT_SEEN = 30 # if device wasn't spotted for this time period, the measure is set to inf
|
||||
|
||||
def __init__(self, trackers: List[str]):
|
||||
self.trackers = trackers
|
||||
self.feature_vecs_per_device = defaultdict(lambda: [self.FAR_AWAY_FEATURE_VALUE] * len(trackers))
|
||||
self.last_measurements = deque()
|
||||
self.tracker_name_to_idx = {name: i for i, name in enumerate(trackers)}
|
||||
self.start_time = None
|
||||
|
||||
@staticmethod
|
||||
def _get_feature_value(rssi, tx_power):
|
||||
"""Transforms rssi and tx power into a value between 0 and 1, where 0 is close and 1 is far away"""
|
||||
MIN_RSSI = -90
|
||||
MAX_TRANSFORMED_RSSI = 40
|
||||
v = tx_power - rssi - MAX_TRANSFORMED_RSSI
|
||||
if v < 0:
|
||||
v = 0
|
||||
return v / (-MIN_RSSI)
|
||||
|
||||
def add_measurement(self, new_measurement: BtleDeviceMeasurement):
|
||||
if self.start_time is None:
|
||||
self.start_time = new_measurement.time
|
||||
|
||||
self.last_measurements.append(new_measurement)
|
||||
while len(self.last_measurements) > 0 and new_measurement.time - self.last_measurements[0].time > self.TIME_TO_DELETE_IF_NOT_SEEN:
|
||||
self.last_measurements.popleft()
|
||||
|
||||
feature_vec = [self.FAR_AWAY_FEATURE_VALUE] * len(self.trackers)
|
||||
for m in self.last_measurements:
|
||||
if m.device == new_measurement.device:
|
||||
tracker_idx = self.tracker_name_to_idx[m.tracker]
|
||||
feature_vec[tracker_idx] = self._get_feature_value(m.rssi, m.tx_power)
|
||||
return feature_vec if new_measurement.time - self.start_time > self.MIN_TIME_UNTIL_PREDICTION else None
|
||||
|
||||
|
||||
|
||||
def training_data_from_df(df: pd.DataFrame, device_to_train: str):
|
||||
"""Returns a feature matrix (num_measurement, num_trackers) and a label vector (both numeric) to be used in scikit learn"""
|
||||
trackers = list(df["tracker"].cat.categories)
|
||||
idx_to_room = dict(enumerate(df["known_room"].cat.categories))
|
||||
room_to_idx = {v: k for k, v in idx_to_room.items()}
|
||||
|
||||
last_known_room = None
|
||||
|
||||
features = []
|
||||
labels = []
|
||||
|
||||
feature_accumulator = RunningFeatureVector(trackers)
|
||||
|
||||
# Feature vectors - rssi column for each room
|
||||
for i, row in df.iterrows():
|
||||
time, device, tracker, rssi, tx_power, known_room = row
|
||||
m = BtleDeviceMeasurement(time, device, tracker, rssi, tx_power)
|
||||
if device != device_to_train:
|
||||
continue
|
||||
if last_known_room != known_room:
|
||||
feature_accumulator = RunningFeatureVector(trackers) # reset for new room
|
||||
last_known_room = known_room
|
||||
feature_vec = feature_accumulator.add_measurement(m)
|
||||
if feature_vec is not None:
|
||||
features.append(feature_vec)
|
||||
labels.append(room_to_idx[known_room])
|
||||
|
||||
return np.array(features), np.array(labels)
|
||||
|
||||
|
||||
def load_measurements_from_csv(csv_file: Path) -> pd.DataFrame:
|
||||
"""Load csv with training data into dataframe"""
|
||||
|
||||
def cleanup_column_name(col_name: str):
|
||||
return col_name.replace("#", "").strip()
|
||||
|
||||
df = pd.read_csv(str(csv_file))
|
||||
|
||||
# String cleanup in column names and room names
|
||||
df = df.rename(columns=cleanup_column_name)
|
||||
df.map(lambda x: x.strip() if isinstance(x, str) else x)
|
||||
|
||||
df["tracker"] = df["tracker"].astype("category")
|
||||
df["known_room"] = df["known_room"].astype("category")
|
||||
df['device'] = df['device'].astype("category")
|
||||
|
||||
return df
|
||||
|
||||
|
||||
async def send_discovery_messages(mqtt_client, device_names):
|
||||
for device_name in device_names:
|
||||
topic = f"homeassistant/sensor/my_btmonitor/{device_name}/config"
|
||||
msg = {
|
||||
"name": device_name,
|
||||
"state_topic": f"my_btmonitor/ml/{device_name}",
|
||||
"expire_after": 30,
|
||||
"unique_id": device_name,
|
||||
}
|
||||
await mqtt_client.publish(topic, json.dumps(msg).encode(), retain=True)
|
||||
|
||||
|
||||
|
||||
async def async_main(
|
||||
mqtt_info: MqttInfo,
|
||||
trackers: List[str],
|
||||
devices: List[str],
|
||||
classifier,
|
||||
device_decoder: DeviceDecoder,
|
||||
training_data_logger: KnownRoomCsvLogger,
|
||||
):
|
||||
current_rooms = defaultdict(lambda: "unknown")
|
||||
feature_accumulator = RunningFeatureVector(trackers)
|
||||
async with aiomqtt.Client(
|
||||
hostname=mqtt_info.server, username=mqtt_info.username, password=mqtt_info.password
|
||||
) as client:
|
||||
await send_discovery_messages(client, devices)
|
||||
await client.subscribe("my_btmonitor/#")
|
||||
async for message in client.messages:
|
||||
current_time = time()
|
||||
topic = message.topic
|
||||
if topic.value == "my_btmonitor/known_room":
|
||||
training_data_logger.update_known_room(message.payload.decode())
|
||||
else:
|
||||
splitted_topic = message.topic.value.split("/")
|
||||
if splitted_topic[0] == "my_btmonitor" and splitted_topic[1] == "raw_measurements":
|
||||
msg_json = json.loads(message.payload)
|
||||
measurement = BtleMeasurement(
|
||||
time=current_time,
|
||||
tracker=splitted_topic[2],
|
||||
address=msg_json["address"],
|
||||
rssi=msg_json["rssi"],
|
||||
tx_power=msg_json.get("tx_power", 0),
|
||||
)
|
||||
logging.debug(f"Got Measurement {measurement}")
|
||||
m = device_decoder(measurement)
|
||||
if m is not None:
|
||||
logging.info(f"Decoded Measurement {m}")
|
||||
training_data_logger.report_measure(m)
|
||||
feature_vec =feature_accumulator.add_measurement(m)
|
||||
if feature_vec:
|
||||
feature_str={tracker : value for tracker, value in zip(trackers, feature_vec)}
|
||||
logging.info(f"Features: {feature_str}")
|
||||
if feature_vec is not None and classifier is not None:
|
||||
room = classifier(m.device, feature_vec)
|
||||
if room != current_rooms[m.device]:
|
||||
logging.info(f"{m.device} moved room {current_rooms[m.device]} to {room}")
|
||||
current_rooms[m.device] = room
|
||||
await client.publish(f"my_btmonitor/ml/{m.device}", room.encode())
|
||||
|
||||
async def async_main_with_restart(
|
||||
mqtt_info: MqttInfo,
|
||||
trackers: List[str],
|
||||
devices: List[str],
|
||||
classifier,
|
||||
device_decoder: DeviceDecoder,
|
||||
training_data_logger: KnownRoomCsvLogger,
|
||||
):
|
||||
while True:
|
||||
try:
|
||||
await async_main(mqtt_info, trackers, devices, classifier, device_decoder, training_data_logger)
|
||||
except Exception as e:
|
||||
print(e)
|
||||
print("restarting...")
|
||||
|
||||
def get_classification_func(training_df: pd.DataFrame, log_classifier_scores=True):
|
||||
devices_to_track = list(training_df["device"].unique())
|
||||
classifiers = {}
|
||||
rooms = list(training_df["known_room"].dtype.categories)
|
||||
for device_to_track in devices_to_track:
|
||||
features, labels = training_data_from_df(training_df, device_to_track)
|
||||
clf = svm.SVC(kernel="rbf")
|
||||
logging.info(f"Computing cross validation score for {device_to_track}")
|
||||
if log_classifier_scores:
|
||||
scores = cross_val_score(clf, features, labels, cv=5)
|
||||
logging.info(" %0.2f accuracy with a standard deviation of %0.2f" % (scores.mean(), scores.std()))
|
||||
|
||||
logging.info(f"Training SVM classifier for {device_to_track}")
|
||||
clf.fit(features, labels)
|
||||
classifiers[device_to_track] = clf
|
||||
|
||||
def classify(device_name, feature_vec):
|
||||
room_idx = classifiers[device_name].predict([feature_vec])[0]
|
||||
return rooms[room_idx]
|
||||
|
||||
return classify
|
||||
|
||||
|
||||
if __name__ == "__main__":
|
||||
mqtt_info = MqttInfo(server="homeassistant.fritz.box", username="my_btmonitor", password="8aBIAC14jaKKbla")
|
||||
# Dict with bt addresses as strings to device name
|
||||
address_to_name = {}
|
||||
# Devices with random addresses - need irk key
|
||||
irk_to_devicename = {
|
||||
"aa67542b82c0e05d65c27fb7e313aba5": "martins_apple_watch",
|
||||
"840e3892644c1ebd1594a9069c14ce0d": "martins_iphone",
|
||||
}
|
||||
script_path = os.path.dirname(os.path.realpath(__file__))
|
||||
data_file = Path(script_path) / Path("training_data.csv")
|
||||
training_df = load_measurements_from_csv(data_file)
|
||||
classification_func = get_classification_func(training_df)
|
||||
training_data_logger = KnownRoomCsvLogger(data_file)
|
||||
device_decoder = DeviceDecoder(irk_to_devicename, address_to_name)
|
||||
trackers = list(training_df["tracker"].cat.categories)
|
||||
devices = list(training_df['device'].cat.categories)
|
||||
asyncio.run(async_main_with_restart(mqtt_info, trackers, devices, classification_func, device_decoder, training_data_logger))
|
||||
27855
roles/bluetooth_monitor/other/collected_backup.csv
Normal file
27855
roles/bluetooth_monitor/other/collected_backup.csv
Normal file
File diff suppressed because it is too large
Load Diff
7
roles/bluetooth_monitor/other/requirements.txt
Normal file
7
roles/bluetooth_monitor/other/requirements.txt
Normal file
@@ -0,0 +1,7 @@
|
||||
aiomqtt==2.0.0
|
||||
numpy==1.26.4
|
||||
pandas==2.2.1
|
||||
pycryptodome==3.20.0
|
||||
scikit-learn==1.4.1.post1
|
||||
scipy==1.12.0
|
||||
typing_extensions==4.10.0
|
||||
27855
roles/bluetooth_monitor/other/training_data.csv
Normal file
27855
roles/bluetooth_monitor/other/training_data.csv
Normal file
File diff suppressed because it is too large
Load Diff
33
roles/bluetooth_monitor/tasks/main.yml
Normal file
33
roles/bluetooth_monitor/tasks/main.yml
Normal file
@@ -0,0 +1,33 @@
|
||||
---
|
||||
- name: Apt install bluez, firmware and Python requirements
|
||||
ansible.builtin.apt:
|
||||
name:
|
||||
- bluez
|
||||
- bluez-firmware
|
||||
- firmware-realtek
|
||||
- firmware-realtek-rtl8723cs-bt
|
||||
- python3-pycryptodome
|
||||
- python3-bleak
|
||||
- python3-asyncio-mqtt
|
||||
- python3-numpy
|
||||
- name: Copy monitor script
|
||||
ansible.builtin.template:
|
||||
src: my_btmonitor.py
|
||||
dest: /usr/bin/my_btmonitor
|
||||
owner: root
|
||||
mode: u+rwx
|
||||
- name: Install systemd service file
|
||||
ansible.builtin.copy:
|
||||
src: my_btmonitor.service
|
||||
dest: /etc/systemd/system/
|
||||
- name: Add script to autostart and start now
|
||||
ansible.builtin.systemd:
|
||||
name: my_btmonitor
|
||||
state: restarted
|
||||
enabled: "yes"
|
||||
daemon_reload: "yes"
|
||||
# - name: Add to sysdweb
|
||||
# include_role:
|
||||
# name: pi_sysdweb
|
||||
# vars:
|
||||
# sysdweb_name: my_btmonitor
|
||||
104
roles/bluetooth_monitor/templates/my_btmonitor.py
Normal file
104
roles/bluetooth_monitor/templates/my_btmonitor.py
Normal file
@@ -0,0 +1,104 @@
|
||||
#!/usr/bin/env python3
|
||||
import asyncio
|
||||
from bleak import BleakScanner
|
||||
from bleak.assigned_numbers import AdvertisementDataType
|
||||
from bleak.backends.bluezdbus.advertisement_monitor import OrPattern
|
||||
from bleak.backends.bluezdbus.scanner import BlueZScannerArgs
|
||||
from functools import partial
|
||||
import asyncio_mqtt
|
||||
import json
|
||||
from datetime import datetime
|
||||
import os
|
||||
import time
|
||||
import subprocess
|
||||
|
||||
# ------------------- Config ----------------------------------------------------------------
|
||||
|
||||
config = {
|
||||
"mqtt": {
|
||||
"hostname": "homeassistant.fritz.box",
|
||||
"username": "{{my_btmonitor_mqtt_username}}",
|
||||
"password": "{{my_btmonitor_mqtt_password}}",
|
||||
"room": "{{sensor_room_name_ascii}}"
|
||||
},
|
||||
"watchdog_seconds": {{my_bt_monitor_watchdog_seconds | default(None)}},
|
||||
"restart_ble_interface": {{my_btmonitor_restart_ble_interface | default(None)}},
|
||||
}
|
||||
|
||||
stop_event = asyncio.Event()
|
||||
time_last_package_received = datetime.now()
|
||||
|
||||
|
||||
async def on_device_found_callback(mqtt_client, room, device, advertising_data):
|
||||
global time_last_package_received
|
||||
time_last_package_received = datetime.now()
|
||||
rssi = advertising_data.rssi
|
||||
tx_power = advertising_data.tx_power
|
||||
if tx_power is not None and rssi is not None:
|
||||
topic = f"my_btmonitor/raw_measurements/{room}"
|
||||
data = {"address": device.address,
|
||||
"rssi": rssi,
|
||||
"tx_power": tx_power}
|
||||
try:
|
||||
await mqtt_client.publish(topic, json.dumps(data).encode())
|
||||
except Exception:
|
||||
print("Probably mqtt isn't running - exit whole script and let systemd restart it")
|
||||
exit(1)
|
||||
|
||||
|
||||
async def watchdog():
|
||||
global time_last_package_received
|
||||
timeout = config["watchdog_seconds"]
|
||||
if not timeout or timeout <= 0:
|
||||
return
|
||||
while True:
|
||||
restart = (datetime.now() - time_last_package_received).seconds > timeout
|
||||
if restart:
|
||||
stop_event.set()
|
||||
await asyncio.sleep(60)
|
||||
|
||||
|
||||
async def ble_scan():
|
||||
mqtt_conf = config['mqtt']
|
||||
while True:
|
||||
try:
|
||||
async with asyncio_mqtt.Client(hostname=mqtt_conf["hostname"],
|
||||
username=mqtt_conf["username"],
|
||||
password=mqtt_conf['password']) as mqtt_client:
|
||||
cb = partial(on_device_found_callback, mqtt_client, mqtt_conf['room'])
|
||||
active_scan = True
|
||||
if active_scan:
|
||||
async with BleakScanner(cb) as scanner:
|
||||
await stop_event.wait()
|
||||
else:
|
||||
# Doesn't work, because of the strange or_patters
|
||||
args = BlueZScannerArgs(
|
||||
or_patterns=[OrPattern(0, AdvertisementDataType.MANUFACTURER_SPECIFIC_DATA, b"\x00\x4c")]
|
||||
)
|
||||
async with BleakScanner(cb, bluez=args, scanning_mode="passive") as scanner:
|
||||
await stop_event.wait()
|
||||
except Exception as e:
|
||||
print("Error", e)
|
||||
try:
|
||||
subprocess.run(["hciconfig", "hci0", "reset"], check=True)
|
||||
except Exception as reset_err:
|
||||
print(f"Reset failed: {reset_err}")
|
||||
await asyncio.sleep(3)
|
||||
|
||||
print("Starting again")
|
||||
|
||||
|
||||
async def main():
|
||||
await asyncio.gather(ble_scan(), watchdog())
|
||||
|
||||
|
||||
if __name__ == "__main__":
|
||||
restart_interface = config["restart_ble_interface"]
|
||||
if restart_interface:
|
||||
print(f"Restarting {restart_interface}")
|
||||
os.system(f"hciconfig {restart_interface} down")
|
||||
time.sleep(3)
|
||||
os.system(f"hciconfig {restart_interface} up")
|
||||
time.sleep(3)
|
||||
print("Done")
|
||||
asyncio.run(main())
|
||||
Reference in New Issue
Block a user