About this page & a note on secrets
The complete source for every script (dissertation Appendix D). Any credentials shown (Wi-Fi, MQTT broker, Firebase service-account file, Telegram/LLM tokens) are placeholders — the real values live only in the private deployment and are never published. Comments have been stripped for brevity.
D.1 Sensor Node Firmware — sensor_estufa.ino
#include <WiFi.h>
#include <MQTT.h>
#include <DHT.h>
#include <Wire.h>
#include <BH1750.h>
#include <ArduinoJson.h>
#define PinoSensorHumidadeSolo 32
#define PowerPIN 4
#define I2C_SDA 25
#define I2C_SCL 26
#define DHTPIN 16
#define TypeDHT DHT11
DHT dht(DHTPIN, TypeDHT);
BH1750 lightMeter;
const int ValorSeco = 3425;
const int ValorMolhado = 1635;
#define TEMPO_DORMIR 300
#define TEMPO_RETRY 60
#define WIFI_TIMEOUT_MS 20000
#define MQTT_TIMEOUT_MS 10000
#define MQTT_PUBACK_TIMEOUT 3000
#define DHT_MAX_TENTATIVAS 5
#define DHT_DELAY_MS 2000
#define ERRO_NENHUM 0
#define ERRO_DHT 1
#define ERRO_WIFI 2
#define ERRO_MQTT 3
const char* ssid = "<your-wifi-ssid>";
const char* password = "<your-wifi-password>";
const char* mqtt_server = "<broker-ip>";
const char* mqtt_user = "<broker-user>";
const char* mqtt_pass = "<broker-password>";
WiFiClient espClient;
MQTTClient client(512);
RTC_DATA_ATTR uint32_t contadorLeituras = 0;
RTC_DATA_ATTR uint16_t ciclosFalhados = 0;
RTC_DATA_ATTR uint8_t ultimoErro = ERRO_NENHUM;
const char* nomeErro(uint8_t e) {
switch (e) {
case ERRO_DHT: return "DHT_NAN";
case ERRO_WIFI: return "WIFI_TIMEOUT";
case ERRO_MQTT: return "MQTT_TIMEOUT";
default: return "NENHUM";
}
}
void entrarDeepSleep(uint32_t segundos = TEMPO_DORMIR) {
client.disconnect();
delay(200);
digitalWrite(PowerPIN, LOW);
WiFi.disconnect(true);
delay(100);
Serial.printf("A entrar em modo Deep Sleep (%u s)...\n\n", segundos);
Serial.flush();
esp_sleep_enable_timer_wakeup((uint64_t)segundos * 1000000ULL);
esp_deep_sleep_start();
}
void abortarCiclo(uint8_t motivo) {
ultimoErro = motivo;
ciclosFalhados++;
Serial.printf("Ciclo abortado: %s (falhas acumuladas: %u)\n",
nomeErro(motivo), ciclosFalhados);
entrarDeepSleep(TEMPO_RETRY);
}
bool lerDHT(float &temperatura, float &humidade) {
for (int i = 1; i <= DHT_MAX_TENTATIVAS; i++) {
humidade = dht.readHumidity();
temperatura = dht.readTemperature();
if (!isnan(humidade) && !isnan(temperatura)) return true;
Serial.printf("DHT: tentativa %d/%d falhou (NaN), a repetir...\n",
i, DHT_MAX_TENTATIVAS);
delay(DHT_DELAY_MS);
}
return false;
}
bool ligarWiFi() {
Serial.print("A ligar ao Wi-Fi...");
WiFi.persistent(false);
WiFi.mode(WIFI_STA);
WiFi.disconnect(true);
delay(100);
WiFi.begin(ssid, password);
unsigned long inicio = millis();
while (WiFi.status() != WL_CONNECTED) {
if (millis() - inicio > WIFI_TIMEOUT_MS) {
Serial.printf("\nWi-Fi: TIMEOUT (status=%d, RSSI=%d dBm)\n",
WiFi.status(), WiFi.RSSI());
return false;
}
delay(500);
Serial.print(".");
}
Serial.printf("\nWi-Fi Ligado! IP=%s RSSI=%d dBm\n",
WiFi.localIP().toString().c_str(), WiFi.RSSI());
return true;
}
bool ligarMQTT() {
client.begin(mqtt_server, 1883, espClient);
Serial.print("A ligar ao MQTT...");
unsigned long inicio = millis();
while (!client.connected()) {
if (millis() - inicio > MQTT_TIMEOUT_MS) {
Serial.println("\nMQTT: timeout atingido.");
return false;
}
if (client.connect("ESP32_Estufa", mqtt_user, mqtt_pass)) {
Serial.println("\nMQTT Ligado!");
return true;
}
Serial.print(".");
delay(2000);
}
return true;
}
void setup() {
Serial.begin(115200);
pinMode(PowerPIN, OUTPUT);
digitalWrite(PowerPIN, HIGH);
delay(1000);
dht.begin();
Wire.begin(I2C_SDA, I2C_SCL);
lightMeter.begin(BH1750::ONE_TIME_HIGH_RES_MODE);
delay(2000);
int valorBrutoSolo = analogRead(PinoSensorHumidadeSolo);
int percentagemSolo = map(valorBrutoSolo, ValorSeco, ValorMolhado, 0, 100);
percentagemSolo = constrain(percentagemSolo, 0, 100);
float humidadeAr, temperaturaAr;
if (!lerDHT(temperaturaAr, humidadeAr)) {
abortarCiclo(ERRO_DHT);
return;
}
float nivelLuz = lightMeter.readLightLevel();
if (nivelLuz < 0) {
Serial.println("Aviso: BH1750 sem leitura. A enviar luminosidade = 0.");
nivelLuz = 0;
}
if (!ligarWiFi()) {
abortarCiclo(ERRO_WIFI);
return;
}
delay(500);
if (!ligarMQTT()) {
abortarCiclo(ERRO_MQTT);
return;
}
contadorLeituras++;
JsonDocument doc;
doc["id"] = "esp32_1_" + String(contadorLeituras);
doc["sensor"] = 1;
doc["humidade_solo"] = percentagemSolo;
doc["temp_ar"] = temperaturaAr;
doc["humidade_ar"] = humidadeAr;
doc["luminosidade"] = nivelLuz;
doc["falhas"] = ciclosFalhados;
doc["ultimo_erro"] = nomeErro(ultimoErro);
char payload[256];
serializeJson(doc, payload);
Serial.print("A enviar mensagem: ");
Serial.println(payload);
bool sucesso = client.publish("sensores/jardim", payload, false, 1);
if (sucesso) {
Serial.println("Sucesso: Mensagem colocada no buffer de envio!");
} else {
Serial.println("Erro: Falha ao colocar a mensagem no buffer!");
}
unsigned long inicioFlush = millis();
while (millis() - inicioFlush < MQTT_PUBACK_TIMEOUT) {
client.loop();
delay(50);
}
ciclosFalhados = 0;
ultimoErro = ERRO_NENHUM;
entrarDeepSleep(TEMPO_DORMIR);
}
void loop() {
}
D.2 Firestore Ingestion — guarda_dados.py
import paho.mqtt.client as mqtt
import json
import firebase_admin
from firebase_admin import credentials
from firebase_admin import firestore
import datetime
import os
import time
MQTT_BROKER = "localhost"
MQTT_TOPIC = "sensores/jardim"
MQTT_USER = "<broker-user>"
MQTT_PASS = "<broker-password>"
BUFFER_FILE = os.environ.get(
"BUFFER_FILE", "/home/sobral03/buffer_pendente.json"
)
MAX_BUFFER_SIZE = 500
MAX_RETRIES = 3
_ft = os.environ.get("FIRESTORE_TIMEOUT")
FIRESTORE_TIMEOUT = float(_ft) if _ft else None
if FIRESTORE_TIMEOUT:
from google.api_core import retry as _g_retry
FS_KWARGS = {"retry": _g_retry.Retry(deadline=FIRESTORE_TIMEOUT),
"timeout": FIRESTORE_TIMEOUT}
else:
FS_KWARGS = {}
DEDUP_WINDOW_S = 120
ids_vistos = {}
try:
cred = credentials.Certificate(
'<project-id>-firebase-adminsdk-<key-id>.json'
)
firebase_admin.initialize_app(cred)
db = firestore.client()
print("[INFO] Firebase ligado com sucesso.")
except Exception as e:
print(f"[FATAL] Erro ao ligar ao Firebase: {e}")
exit(1)
def carregar_buffer():
if not os.path.exists(BUFFER_FILE):
return []
try:
with open(BUFFER_FILE, "r") as f:
dados = json.load(f)
if isinstance(dados, list):
return dados
except (json.JSONDecodeError, IOError) as e:
print(f"[WARN] Buffer corrompido, a descartar: {e}")
return []
def guardar_buffer(buffer):
tmp = f"{BUFFER_FILE}.tmp"
try:
with open(tmp, "w") as f:
json.dump(buffer, f, default=str)
os.replace(tmp, BUFFER_FILE)
except (IOError, OSError, TypeError) as e:
print(f"[ERROR] Falha ao guardar buffer local: {e}")
if os.path.exists(tmp):
try:
os.remove(tmp)
except OSError:
pass
def adicionar_ao_buffer(dado):
buffer = carregar_buffer()
buffer.append(dado)
if len(buffer) > MAX_BUFFER_SIZE:
descartados = len(buffer) - MAX_BUFFER_SIZE
buffer = buffer[-MAX_BUFFER_SIZE:]
print(f"[WARN] Buffer cheio. {descartados} leituras antigas descartadas.")
guardar_buffer(buffer)
print(f"[BUFFER] Leitura guardada localmente. Total pendentes: {len(buffer)}")
def flush_buffer():
buffer = carregar_buffer()
if not buffer:
return
print(f"[BUFFER] A tentar enviar {len(buffer)} leituras pendentes...")
falhas = []
for dado in buffer:
try:
envio = dict(dado)
if isinstance(envio.get("data_hora"), str):
envio["data_hora"] = datetime.datetime.fromisoformat(envio["data_hora"])
db.collection('estufa').add(envio, **FS_KWARGS)
except Exception:
falhas.append(dado)
if falhas:
guardar_buffer(falhas)
print(f"[BUFFER] {len(buffer) - len(falhas)} enviadas, "
f"{len(falhas)} ainda pendentes.")
else:
if os.path.exists(BUFFER_FILE):
os.remove(BUFFER_FILE)
print(f"[BUFFER] Todas as {len(buffer)} leituras pendentes enviadas.")
def escrever_firestore(dados):
for tentativa in range(1, MAX_RETRIES + 1):
try:
db.collection('estufa').add(dados, **FS_KWARGS)
print(f"[OK] Dados guardados: Sensor {dados.get('sensor')} | "
f"Temp: {dados.get('temp_ar')} C | "
f"Hum. Solo: {dados.get('humidade_solo')}%")
return True
except Exception as e:
print(f"[ERROR] Tentativa {tentativa}/{MAX_RETRIES} falhou: {e}")
if tentativa < MAX_RETRIES:
time.sleep(2)
dados_serializavel = dict(dados)
if isinstance(dados_serializavel.get("data_hora"), datetime.datetime):
dados_serializavel["data_hora"] = dados_serializavel["data_hora"].isoformat()
adicionar_ao_buffer(dados_serializavel)
return False
def on_connect(client, userdata, flags, rc):
if rc == 0:
print(f"[INFO] Ligado ao MQTT. A ouvir o topico '{MQTT_TOPIC}'...")
client.subscribe(MQTT_TOPIC, qos=1)
else:
print(f"[ERROR] Falha ao ligar ao MQTT. Codigo: {rc}")
def on_message(client, userdata, msg):
try:
payload = msg.payload.decode('utf-8')
dados = json.loads(payload)
for campo in ("falhas", "ultimo_erro"):
dados.pop(campo, None)
msg_id = dados.get("id")
agora = time.monotonic()
visto = ids_vistos.get(msg_id)
if msg_id and visto is not None and agora - visto <= DEDUP_WINDOW_S:
print(f"[DEDUP] {msg_id} reentregue ha {agora - visto:.0f}s "
f"(<{DEDUP_WINDOW_S}s), a ignorar.")
return
dados['data_hora'] = datetime.datetime.now(datetime.timezone.utc)
if escrever_firestore(dados) and msg_id:
ids_vistos[msg_id] = agora
for k in [k for k, t in ids_vistos.items()
if agora - t > DEDUP_WINDOW_S]:
del ids_vistos[k]
flush_buffer()
except json.JSONDecodeError as e:
print(f"[ERROR] JSON invalido: {e}")
except Exception as e:
print(f"[ERROR] Erro inesperado: {e}")
flush_buffer()
cliente_mqtt = mqtt.Client()
cliente_mqtt.username_pw_set(MQTT_USER, MQTT_PASS)
cliente_mqtt.on_connect = on_connect
cliente_mqtt.on_message = on_message
cliente_mqtt.connect(MQTT_BROKER, 1883, 60)
print("[INFO] A arrancar... (Pressiona Ctrl+C para sair)")
cliente_mqtt.loop_forever()
D.3 Live Monitor — monitor.py
import datetime
import json
import os
import threading
import time
import urllib.parse
import urllib.request
import paho.mqtt.client as mqtt
MQTT_BROKER = os.environ.get("MQTT_BROKER", "localhost")
MQTT_PORT = int(os.environ.get("MQTT_PORT", "1883"))
MQTT_TOPIC = os.environ.get("MQTT_TOPIC", "sensores/jardim")
MQTT_USER = os.environ.get("MQTT_USER", "estufa")
MQTT_PASS = os.environ.get("MQTT_PASS", "admin")
SENSOR_FIELD = os.environ.get("SENSOR_FIELD", "sensor")
MOISTURE_FIELD = os.environ.get("MOISTURE_FIELD", "humidade_solo")
TEMP_FIELD = os.environ.get("TEMP_FIELD", "temp_ar")
HUMIDITY_FIELD = os.environ.get("HUMIDITY_FIELD", "humidade_ar")
LUX_FIELD = os.environ.get("LUX_FIELD", "luminosidade")
FALHAS_FIELD = os.environ.get("FALHAS_FIELD", "falhas")
ERRO_FIELD = os.environ.get("ERRO_FIELD", "ultimo_erro")
STALE_SECONDS = int(os.environ.get("ESTUFA_STALE_SECONDS", "1800"))
CHECK_SECONDS = int(os.environ.get("ESTUFA_CHECK_SECONDS", "60"))
STATE_FILE = os.environ.get(
"ESTUFA_STATE_FILE", os.path.expanduser("~/.openclaw/monitor/state.json")
)
TG_TOKEN = os.environ.get("TELEGRAM_BOT_TOKEN", "")
TG_CHAT = os.environ.get("TELEGRAM_CHAT_ID", "")
sensors = {}
lock = threading.Lock()
def iso(epoch):
return datetime.datetime.fromtimestamp(epoch, datetime.timezone.utc).isoformat()
def telegram(text):
print(f"[ALERT] {text}", flush=True)
if not TG_TOKEN or not TG_CHAT:
print("[ALERT] (Telegram nao configurado ainda -- so registado no log)", flush=True)
return
try:
url = f"https://api.telegram.org/bot{TG_TOKEN}/sendMessage"
data = urllib.parse.urlencode({"chat_id": TG_CHAT, "text": text}).encode()
urllib.request.urlopen(urllib.request.Request(url, data=data), timeout=10)
except Exception as e:
print(f"[WARN] Falha a enviar Telegram: {e}", flush=True)
def load_state():
try:
with open(STATE_FILE) as f:
saved = json.load(f)
except (OSError, ValueError):
return
for sid, s in saved.get("sensors", {}).items():
try:
last_seen = datetime.datetime.fromisoformat(s["last_seen"]).timestamp()
except (KeyError, TypeError, ValueError):
continue
online = bool(s.get("online", True))
sensors[sid] = {
"last_seen": last_seen,
"online": online,
"alerted": not online,
"moisture": s.get("moisture"),
"temp": s.get("temp"),
"humidity": s.get("humidity"),
"lux": s.get("lux"),
"falhas": s.get("falhas"),
"ultimo_erro": s.get("ultimo_erro"),
}
if sensors:
print(
f"[INFO] Estado recuperado: {len(sensors)} sensor(es) conhecidos: "
f"{', '.join(sensors)}",
flush=True,
)
def write_state():
now = time.time()
payload = {
"updated": iso(now),
"stale_seconds": STALE_SECONDS,
"sensors": {
sid: {
"online": s["online"],
"last_seen": iso(s["last_seen"]),
"age_seconds": round(now - s["last_seen"]),
"moisture": s.get("moisture"),
"temp": s.get("temp"),
"humidity": s.get("humidity"),
"lux": s.get("lux"),
"falhas": s.get("falhas"),
"ultimo_erro": s.get("ultimo_erro"),
}
for sid, s in sensors.items()
},
}
os.makedirs(os.path.dirname(STATE_FILE), exist_ok=True)
tmp = STATE_FILE + ".tmp"
with open(tmp, "w") as f:
json.dump(payload, f, indent=2, ensure_ascii=False)
os.replace(tmp, STATE_FILE)
def motivo_recuperacao(d):
falhas = d.get(FALHAS_FIELD)
if not falhas:
return ""
erro = d.get(ERRO_FIELD, "?")
return f" (tinha falhado {falhas} ciclo(s); ultimo motivo: {erro})"
def on_connect(client, userdata, flags, reason_code, properties=None):
if reason_code == 0:
print(f"[INFO] Ligado ao MQTT, a ouvir '{MQTT_TOPIC}'", flush=True)
client.subscribe(MQTT_TOPIC, qos=1)
else:
print(f"[ERROR] Falha ao ligar ao MQTT: {reason_code}", flush=True)
def on_message(client, userdata, msg):
try:
d = json.loads(msg.payload.decode("utf-8"))
except Exception as e:
print(f"[WARN] payload invalido: {e}", flush=True)
return
sid = str(d.get(SENSOR_FIELD, "desconhecido"))
with lock:
is_new = sid not in sensors
s = sensors.setdefault(sid, {"online": True, "alerted": False})
was_offline = not s.get("online", True)
s["last_seen"] = time.time()
s["online"] = True
s["alerted"] = False
if MOISTURE_FIELD in d:
s["moisture"] = d.get(MOISTURE_FIELD)
if TEMP_FIELD in d:
s["temp"] = d.get(TEMP_FIELD)
if HUMIDITY_FIELD in d:
s["humidity"] = d.get(HUMIDITY_FIELD)
if LUX_FIELD in d:
s["lux"] = d.get(LUX_FIELD)
if FALHAS_FIELD in d:
s["falhas"] = d.get(FALHAS_FIELD)
if ERRO_FIELD in d:
s["ultimo_erro"] = d.get(ERRO_FIELD)
if is_new:
telegram(f"\U0001F331 Novo sensor detetado: {sid} (a monitorizar).")
elif was_offline:
telegram(f"✅ Sensor {sid} voltou a responder.{motivo_recuperacao(d)}")
write_state()
def watchdog():
while True:
time.sleep(CHECK_SECONDS)
with lock:
now = time.time()
for sid, s in sensors.items():
age = now - s["last_seen"]
if age > STALE_SECONDS and s["online"]:
s["online"] = False
if not s["online"] and not s.get("alerted"):
s["alerted"] = True
telegram(
f"⚠️ Sensor {sid} sem resposta ha {int(age // 60)} min."
)
write_state()
def main():
os.makedirs(os.path.dirname(STATE_FILE), exist_ok=True)
load_state()
write_state()
threading.Thread(target=watchdog, daemon=True).start()
client = mqtt.Client(mqtt.CallbackAPIVersion.VERSION2)
client.username_pw_set(MQTT_USER, MQTT_PASS)
client.on_connect = on_connect
client.on_message = on_message
client.connect(MQTT_BROKER, MQTT_PORT, 60)
print("[INFO] Monitor da estufa a arrancar...", flush=True)
client.loop_forever()
if __name__ == "__main__":
main()
D.4 Skill estufa-status — status.py
import datetime
import json
import os
STATE_FILE = os.environ.get(
"ESTUFA_STATE_FILE", os.path.expanduser("~/.openclaw/monitor/state.json")
)
def main():
if not os.path.exists(STATE_FILE):
print("Sem dados: o monitor ainda nao escreveu estado "
"(o servico estufa-monitor esta a correr?).")
return
try:
with open(STATE_FILE) as f:
state = json.load(f)
except (json.JSONDecodeError, IOError) as e:
print(f"Nao consegui ler o estado do monitor: {e}")
return
sensors = state.get("sensors", {})
if not sensors:
print("O monitor esta a correr mas ainda nao recebeu nenhuma leitura "
"de nenhum sensor.")
return
stale_min = state.get("stale_seconds", 1800) // 60
print(f"Estado dos sensores (sem resposta = silencio > {stale_min} min):")
for sid in sorted(sensors):
s = sensors[sid]
age_min = s.get("age_seconds", 0) / 60
state_txt = "OK" if s.get("online") else "SEM RESPOSTA"
parts = []
moisture = s.get("moisture")
parts.append(f"humidade solo {moisture}%" if moisture is not None else "humidade solo: sem leitura")
temp = s.get("temp")
if temp is not None:
parts.append(f"temp ar {temp}C")
humidity = s.get("humidity")
if humidity is not None:
parts.append(f"humidade ar {humidity}%")
lux = s.get("lux")
if lux is not None:
parts.append(f"luz {lux}")
readings = ", ".join(parts)
print(f"- Sensor {sid}: {readings}; ultima msg ha {age_min:.0f} min [{state_txt}]")
updated = state.get("updated")
if updated:
try:
ts = datetime.datetime.fromisoformat(updated)
print(f"(estado atualizado: {ts:%Y-%m-%d %H:%M:%S %Z})")
except ValueError:
pass
if __name__ == "__main__":
main()
D.5 Skill estufa-diag — diag.py
import datetime
import json
import os
import subprocess
import sys
SYSTEM_SERVICES = ["estufa-monitor", "estufa", "guarda_dados"]
USER_SERVICES = ["openclaw-gateway"]
STATE_FILE = os.path.expanduser("~/.openclaw/monitor/state.json")
def run(cmd, timeout=15):
try:
p = subprocess.run(
cmd, capture_output=True, text=True, timeout=timeout
)
out = (p.stdout or "") + (p.stderr or "")
return out.strip() or "(sem output)"
except subprocess.TimeoutExpired:
return f"(timeout ao correr: {' '.join(cmd)})"
except FileNotFoundError:
return f"(comando nao encontrado: {cmd[0]})"
except Exception as e:
return f"(erro a correr {' '.join(cmd)}: {e})"
def section(title):
print(f"\n=== {title} ===")
def pi_health():
section("Saude do Raspberry Pi")
print("Uptime/carga:", run(["uptime"]))
print("Temp CPU:", run(["vcgencmd", "measure_temp"]))
print("Memoria:\n" + run(["free", "-h"]))
print("Disco (/):\n" + run(["df", "-h", "/"]))
def services():
section("Servicos")
for svc in SYSTEM_SERVICES:
active = run(["systemctl", "is-active", svc])
status = run(["systemctl", "status", svc, "--no-pager", "-n", "0"])
head = "\n".join(status.splitlines()[:4])
print(f"- {svc}: [{active}]\n{head}")
for svc in USER_SERVICES:
active = run(["systemctl", "--user", "is-active", svc])
status = run(["systemctl", "--user", "status", svc, "--no-pager", "-n", "0"])
head = "\n".join(status.splitlines()[:4])
print(f"- {svc} (user): [{active}]\n{head}")
def mqtt_broker():
section("Broker MQTT (mosquitto)")
print("mosquitto:", run(["systemctl", "is-active", "mosquitto"]))
ss = run(["ss", "-ltnp"])
lines = [l for l in ss.splitlines() if ":1883" in l]
print("porta 1883:", "\n".join(lines) if lines else "(nao parece estar a escutar em 1883)")
def monitor_state():
section("Estado do monitor (state.json)")
if not os.path.exists(STATE_FILE):
print("state.json nao existe -- o estufa-monitor ainda nao escreveu estado.")
return
try:
with open(STATE_FILE) as f:
st = json.load(f)
except Exception as e:
print(f"Nao consegui ler state.json: {e}")
return
print("Atualizado:", st.get("updated"))
stale = st.get("stale_seconds", 1800)
for sid, s in sorted(st.get("sensors", {}).items()):
age = s.get("age_seconds", 0)
flag = "OK" if s.get("online") else "SEM RESPOSTA"
warn = " <-- ATRASADO" if age > 360 else ""
print(
f"- Sensor {sid}: solo={s.get('moisture')}% temp={s.get('temp')}C "
f"hum_ar={s.get('humidity')}% luz={s.get('lux')} | "
f"ultima ha {age}s [{flag}]{warn}"
)
print(f"(limiar de silencio do monitor: {stale}s)")
def recent_logs(service):
section(f"Logs recentes: {service}")
out = run(["journalctl", "-u", service, "-n", "30", "--no-pager", "--since", "-30min"])
print(out)
def main():
now = datetime.datetime.now(datetime.timezone.utc)
print(f"Diagnostico estufa @ {now:%Y-%m-%d %H:%M:%S %Z}")
pi_health()
services()
mqtt_broker()
monitor_state()
recent_logs("estufa-monitor")
if len(sys.argv) > 1:
svc = sys.argv[1].strip()
if svc and svc != "estufa-monitor":
recent_logs(svc)
if __name__ == "__main__":
main()
D.6 Skill estufa-rega — water.py
import json
import os
import sys
import paho.mqtt.publish as publish
BROKER = os.environ.get("MQTT_BROKER", "localhost")
PORT = int(os.environ.get("MQTT_PORT", "1883"))
TOPIC = os.environ.get("REGA_TOPIC", "rega/comando")
MQTT_USER = os.environ.get("MQTT_USER", "estufa")
MQTT_PASS = os.environ.get("MQTT_PASS", "admin")
DURACAO_MAX_S = 120
INTENSIDADE_MIN = 20
INTENSIDADE_MAX = 100
def enviar(payload: dict) -> None:
publish.single(
TOPIC,
json.dumps(payload),
qos=1,
hostname=BROKER,
port=PORT,
auth={"username": MQTT_USER, "password": MQTT_PASS},
)
if len(sys.argv) > 1 and sys.argv[1].lower() in ("parar", "stop"):
enviar({"acao": "parar"})
print(f"Sent: PARAR -> {TOPIC} (QoS 1).")
sys.exit(0)
duracao = int(sys.argv[1]) if len(sys.argv) > 1 else 30
intensidade = int(sys.argv[2]) if len(sys.argv) > 2 else 80
capped = duracao > DURACAO_MAX_S
duracao = max(1, min(duracao, DURACAO_MAX_S))
intensidade = max(INTENSIDADE_MIN, min(intensidade, INTENSIDADE_MAX))
enviar({"acao": "regar", "intensidade": intensidade, "duracao": duracao})
msg = f"Sent: regar {duracao}s @ {intensidade}% -> {TOPIC} (QoS 1)."
if capped:
msg += f" NOTA: duracao pedida excedia o maximo, limitada a {DURACAO_MAX_S}s."
print(msg)
D.7 Skill estufa-restart — restart.py
import subprocess
import sys
WHITELIST = {
"estufa-monitor": "system",
"guarda_dados": "system",
"estufa": "system",
"openclaw-gateway": "user",
}
def run(cmd, timeout=30):
try:
p = subprocess.run(cmd, capture_output=True, text=True, timeout=timeout)
return p.returncode, ((p.stdout or "") + (p.stderr or "")).strip()
except subprocess.TimeoutExpired:
return 124, f"(timeout ao correr: {' '.join(cmd)})"
except Exception as e:
return 1, f"(erro: {e})"
def main():
if len(sys.argv) != 2:
print("Uso: restart.py SERVICE (servicos permitidos: "
+ ", ".join(WHITELIST) + ")")
return
svc = sys.argv[1].strip()
if svc not in WHITELIST:
print(f"RECUSADO: '{svc}' nao esta na lista branca. "
f"So posso reiniciar: {', '.join(WHITELIST)}.")
return
scope = WHITELIST[svc]
if scope == "system":
restart_cmd = ["sudo", "-n", "systemctl", "restart", svc]
active_cmd = ["systemctl", "is-active", svc]
else:
restart_cmd = ["systemctl", "--user", "restart", svc]
active_cmd = ["systemctl", "--user", "is-active", svc]
print(f"A reiniciar '{svc}'...")
rc, out = run(restart_cmd)
if rc != 0:
print(f"FALHOU o restart de '{svc}' (codigo {rc}): {out}")
return
_, active = run(active_cmd)
if active == "active":
print(f"OK: '{svc}' reiniciado e esta [active].")
else:
print(f"ATENCAO: '{svc}' foi reiniciado mas o estado agora e "
f"[{active}] -- verifica os logs com a skill estufa-diag.")
if __name__ == "__main__":
main()
D.8 SQLite Logger — regista_sqlite.py
import datetime
import json
import os
import sqlite3
import time
import paho.mqtt.client as mqtt
MQTT_BROKER = os.environ.get("MQTT_BROKER", "localhost")
MQTT_PORT = int(os.environ.get("MQTT_PORT", "1883"))
MQTT_USER = os.environ.get("MQTT_USER", "estufa")
MQTT_PASS = os.environ.get("MQTT_PASS", "admin")
TOPICO_SENSOR = "sensores/jardim"
TOPICO_DECISAO = "rega/decisao"
TOPICO_ESTADO = "rega/estado"
DB_FILE = os.environ.get("ESTUFA_DB", "/home/pi/estufa.db")
DEDUP_WINDOW_S = 120
ids_vistos = {}
SCHEMA =
db = sqlite3.connect(DB_FILE, check_same_thread=False)
db.executescript(SCHEMA)
db.commit()
def agora_iso():
return datetime.datetime.now(datetime.timezone.utc).isoformat()
def on_connect(client, userdata, flags, rc):
if rc == 0:
print(f"[INFO] Ligado ao MQTT. A registar em {DB_FILE}")
for t in (TOPICO_SENSOR, TOPICO_DECISAO, TOPICO_ESTADO):
client.subscribe(t, qos=1)
else:
print(f"[ERROR] Falha ao ligar ao MQTT. Codigo: {rc}")
def on_message(client, userdata, msg):
try:
dados = json.loads(msg.payload.decode("utf-8"))
except (UnicodeDecodeError, json.JSONDecodeError) as e:
print(f"[WARN] payload invalido em {msg.topic}: {e}")
return
ts = agora_iso()
try:
if msg.topic == TOPICO_SENSOR:
msg_id = dados.get("id")
t = time.monotonic()
visto = ids_vistos.get(msg_id)
if msg_id and visto is not None and t - visto <= DEDUP_WINDOW_S:
print(f"[DEDUP] {msg_id} reentregue, a ignorar.")
return
db.execute(
"INSERT INTO leituras (id_leitura, sensor, humidade_solo,"
" temp_ar, humidade_ar, luminosidade, falhas, ultimo_erro,"
" data_hora) VALUES (?,?,?,?,?,?,?,?,?)",
(msg_id, dados.get("sensor"), dados.get("humidade_solo"),
dados.get("temp_ar"), dados.get("humidade_ar"),
dados.get("luminosidade"), dados.get("falhas"),
dados.get("ultimo_erro"), ts),
)
db.commit()
if msg_id:
ids_vistos[msg_id] = t
for k in [k for k, v in ids_vistos.items()
if t - v > DEDUP_WINDOW_S]:
del ids_vistos[k]
print(f"[OK] leitura sensor {dados.get('sensor')}: "
f"solo {dados.get('humidade_solo')}%")
elif msg.topic == TOPICO_DECISAO:
db.execute("INSERT INTO decisoes (payload, data_hora) VALUES (?,?)",
(json.dumps(dados), ts))
db.commit()
elif msg.topic == TOPICO_ESTADO:
db.execute(
"INSERT INTO eventos_rega (estado, intensidade, motivo,"
" payload, data_hora) VALUES (?,?,?,?,?)",
(dados.get("estado"), dados.get("intensidade"),
dados.get("motivo"), json.dumps(dados), ts),
)
db.commit()
except sqlite3.Error as e:
print(f"[ERROR] SQLite: {e}")
try:
cliente = mqtt.Client(mqtt.CallbackAPIVersion.VERSION1)
except (AttributeError, TypeError):
cliente = mqtt.Client()
cliente.username_pw_set(MQTT_USER, MQTT_PASS)
cliente.on_connect = on_connect
cliente.on_message = on_message
cliente.connect(MQTT_BROKER, MQTT_PORT, 60)
print("[INFO] A arrancar... (Ctrl+C para sair)")
cliente.loop_forever()
D.9 Random Forest Training — treinar_rf.py
import sqlite3
import sys
import numpy as np
import pandas as pd
from sklearn.ensemble import RandomForestRegressor
from sklearn.metrics import mean_absolute_error, r2_score
from sklearn.model_selection import train_test_split
DB_FILE = "/home/pi/estufa.db"
HORIZONTE = int(sys.argv[1]) if len(sys.argv) > 1 else 12
con = sqlite3.connect(DB_FILE)
df = pd.read_sql_query(
"SELECT data_hora, sensor, humidade_solo, temp_ar, humidade_ar,"
" luminosidade FROM leituras WHERE sensor = 1 ORDER BY data_hora",
con, parse_dates=["data_hora"],
)
con.close()
if len(df) < HORIZONTE + 50:
sys.exit(f"Dados insuficientes: {len(df)} leituras "
f"(precisas de pelo menos {HORIZONTE + 50}). Deixa o "
f"regista_sqlite.py a acumular mais tempo.")
hora = df["data_hora"].dt.hour + df["data_hora"].dt.minute / 60.0
df["hora_sin"] = np.sin(2 * np.pi * hora / 24)
df["hora_cos"] = np.cos(2 * np.pi * hora / 24)
df["delta_solo"] = df["humidade_solo"].diff(6)
df["alvo"] = df["humidade_solo"].shift(-HORIZONTE)
df = df.dropna()
FEATURES = ["humidade_solo", "temp_ar", "humidade_ar", "luminosidade",
"hora_sin", "hora_cos", "delta_solo"]
X, y = df[FEATURES], df["alvo"]
X_tr, X_te, y_tr, y_te = train_test_split(X, y, test_size=0.2, shuffle=False)
modelo = RandomForestRegressor(n_estimators=300, min_samples_leaf=3,
random_state=42, n_jobs=-1)
modelo.fit(X_tr, y_tr)
pred = modelo.predict(X_te)
print(f"Leituras usadas: {len(df)} | horizonte: {HORIZONTE} (~{HORIZONTE*5} min)")
print(f"MAE: {mean_absolute_error(y_te, pred):.2f} pontos % de humidade do solo")
print(f"R^2: {r2_score(y_te, pred):.3f}")
print("\nImportancia das features:")
for nome, imp in sorted(zip(FEATURES, modelo.feature_importances_),
key=lambda p: -p[1]):
print(f" {nome:15s} {imp:.3f}")
try:
import joblib
joblib.dump(modelo, "/home/pi/modelo_rf.joblib")
print("\nModelo guardado em /home/pi/modelo_rf.joblib")
except ImportError:
print("\n(joblib indisponivel — modelo nao guardado)")