Skip to content

Kernkomponenten

LXMFBot

Die Hauptklasse des Bots. Sie übernimmt das Routing von Nachrichten, die Befehlsverarbeitung und die Verwaltung des Bot-Lebenszyklus.

from lxmfy import LXMFBot

bot = LXMFBot(
    name="MyBot",
    announce=600,
    announce_immediately=True,
    admins=set(),
    hot_reloading=False,
    rate_limit=5,
    cooldown=60,
    max_warnings=3,
    warning_timeout=300,
    command_prefix="/",
    cogs_dir="cogs",
    cogs_enabled=True,
    permissions_enabled=False,
    storage_type="json", # "json", "sqlite" oder "memory"
    storage_path="data",
    first_message_enabled=True,
    event_logging_enabled=True,
    max_logged_events=1000,
    event_middleware_enabled=True,
    announce_enabled=True,
    signature_verification_enabled=False,
    require_message_signatures=False,
    identity_pinning_enabled=False,
    message_persistence_enabled=True,
    dynamic_cogs_enabled=True,
    external_cogs_enabled=True,
    external_cogs_sandbox_enabled=True,
    external_cogs_sandbox_type="auto",  # "auto", "landlock", "bwrap", "firejail", "none"
    external_cogs_timeout=30,
    landlock_enabled=True,
    nlp_enabled=False,
    nlp_threshold=0.5,
    link_support_enabled=False,
    lxmf_commands_enabled=True,
    message_queue_size=50,
    reticulum_config_dir=None,  # oder LXMFY_RETICULUM_CONFIG_DIR / "~/.reticulum"
    rrc_enabled=False,
    rrc_hubs=[],
    rrc_rooms=[],
    rrc_nick=None,
    rrc_dest_name="rrc.hub",
    rrc_auto_reconnect=True,
    rrc_persist_sessions=True,
)

Wichtige Methoden

  • get_landlock_status(): Verfügbarkeit und Aktivierungsstatus der Landlock-LSM-Sandbox für den Bot-Prozess zurückgeben
  • run(delay=10): Hauptschleife des Bots starten
  • send(destination, message, title="Reply", lxmf_fields=None, stamp_cost=None, opportunistic=None): Nachricht an ein Ziel senden, optional mit eigenen LXMF-Feldern, Überschreibung der Stamp-Kosten und opportunistischem Versand (versucht Direktzustellung, greift bei Konfiguration sofort auf Propagation zurück).
  • send_with_attachment(destination, message, attachment, title="Reply", stamp_cost=None, opportunistic=None): Nachricht mit Anhang senden
  • command(name, description="No description provided", admin_only=False, threaded=False): Dekorator zur Registrierung von Befehlen. Mit threaded=True läuft der Befehls-Callback in einem separaten Thread. Befehle unterstützen typannotierte Argumente für automatische Konvertierung.
  • intent(name, examples): Dekorator zur Registrierung von NLP-Intent-Handlern.
  • nlp.export_model(): Trainierte NLP-Modelldaten exportieren.
  • nlp.import_model(model_data): Zuvor exportierte NLP-Modelldaten importieren.
  • request_link(destination_hash, callback=None, app_name="lxmf", *aspects): RNS-Link zu einem Ziel anfordern. Erlaubt eigenes app_name und aspects (Standard "lxmf" und "delivery").
  • on_link(callback): Handler für eingehende RNS-Links registrieren.
  • load_extension(name): Cog-Erweiterungsmodul anhand des Namens laden (z. B. "cogs.utility").
  • reload_extension(name): Cog-Erweiterungsmodul neu laden.
  • add_cog(cog_instance): Cog-Klasseninstanz zum Bot hinzufügen.
  • remove_cog(cog_name): Cog anhand seines Klassennamens entfernen.
  • on_first_message(): Dekorator zur Verarbeitung erster Nachrichten von Benutzern
  • on_message(): Dekorator zur Verarbeitung aller Nachrichten (wird vor der Befehlsverarbeitung aufgerufen)
  • on_reaction(): Dekorator zur Verarbeitung eingehender Reaktionen. Handler erhalten (sender, reaction), wobei reaction die Schlüssel reaction_to, reaction_emoji und reaction_sender enthält
  • react(destination, message_hash, reaction): Reaktion auf eine Nachricht über das LXMF-Feld FIELD_REACTION senden
  • validate(): Validierungsprüfungen der Bot-Konfiguration ausführen
  • connect_rrc(hub_hash, rooms=None, nick=None, dest_name=None, auto_reconnect=None): Als Client mit einem RRC-Hub verbinden
  • disconnect_rrc(hub_hash=None): Eine oder alle RRC-Hub-Sitzungen trennen
  • on_rrc(callback=None): Dekorator oder Handler-Registrierung für RRC-Ereignisse (handler(event, client, payload))
  • rrc: RRCManager-Instanz für Sitzungen mit mehreren Hubs

Strukturierte Befehle über LXMF-Felder

Bots können Befehle empfangen, die über das LXMF-Feld FIELD_COMMANDS (0x09) gesendet werden, und antworten automatisch mit FIELD_RESULTS (0x0A). Damit sind strukturierte Anfrage/Antwort-Workflows neben normalen Textbefehlen möglich.

Eingehende FIELD_COMMANDS werden geparst und über dieselbe Befehlsregistrierung wie Textbefehle geroutet, mit gemeinsamen Berechtigungsprüfungen, typannotiertem Argument-Parsing, Threading und Middleware.

from lxmfy import LXMFBot, FIELD_COMMANDS, FIELD_RESULTS, pack_result, unpack_commands

bot = LXMFBot(name="FieldBot")

@bot.command(name="status", description="Return bot status")
def status_cmd(ctx):
    # ctx.fields enthält das rohe LXMF-Felder-Dict
    # ctx.request_id wird automatisch gesetzt, wenn der Befehl eine enthielt
    ctx.reply("Bot is online")

# Senden eines strukturierten Befehls von einem anderen LXMF-Client:
# lxm.fields[FIELD_COMMANDS] = {"command": "status", "args": [], "request_id": "abc123"}
# router.handle_outbound(lxm)

# Die Bot-Antwort enthält automatisch FIELD_RESULTS mit Antwort und request_id.

Um die Verarbeitung von Feldbefehlen zu deaktivieren, lxmf_commands_enabled=False in BotConfig setzen.

Reaktionen

Reaktionen werden als LXMF-Feld FIELD_REACTION (0x40) in einer ansonsten leeren Nachricht übertragen. pack_reaction und unpack_reaction erzeugen und parsen dieses Feld.

from lxmfy import pack_reaction, unpack_reaction

# Eine Reaktion auf eine Nachricht senden
bot.react(destination_hash, message_hash_hex, "thumbs up emoji")

# Reaktionen empfangen
@bot.on_reaction()
def on_reaction(sender, reaction):
    # reaction["reaction_to"]  - Hex-Hash der Zielnachricht
    # reaction["reaction_emoji"] - Reaktionstext (bis zu 16 Zeichen)
    # reaction["reaction_sender"] - Absender
    print(f"{sender} reacted {reaction['reaction_emoji']} to {reaction['reaction_to']}")
    return True

Der Reaktionstext ist auf 16 druckbare Zeichen begrenzt. Das Rohfeld bleibt in ctx.fields und msg.fields verfügbar (Kompatibilität).

Speicher

Das Framework bietet drei Speicher-Backends:

JSONStorage

from lxmfy import JSONStorage

storage = JSONStorage("data")

SQLiteStorage

from lxmfy import SQLiteStorage

storage = SQLiteStorage("data/bot.db")

MemoryStorage

from lxmfy.storage import MemoryStorage

storage = MemoryStorage() # Vollständig im Arbeitsspeicher

Befehle

Befehlsregistrierung und -verarbeitung:

@bot.command(name="hello", description="Says hello")
def hello(ctx):
    ctx.reply(f"Hello {ctx.sender}!")

Typannotierte Argumente

Befehle parsen und konvertieren Argumente automatisch anhand der Typannotationen in der Callback-Funktion.

@bot.command(name="add", description="Adds two numbers")
def add(ctx, a: int, b: int):
    result = a + b
    ctx.reply(f"The result is {result}")

Hilfesystem

Das Framework enthält einen interaktiven Hilfe-Generator, der kategorisierte Hilfemenüs auf Basis der Cog- und Command-Metadaten erzeugt.

# Der help-Befehl wird automatisch registriert.
# Benutzer können '/help' oder '/help <command>' verwenden

Befehle in separaten Threads

Für lang laufende oder blockierende Operationen, die nicht direkt mit dem Reticulum Network Stack interagieren, können Befehle in einem separaten Thread laufen, damit der Bot reaktionsfähig bleibt.

import time

@bot.command(name="long_task", description="Performs a long-running task in a separate thread", threaded=True)
def long_task_command(ctx):
    ctx.reply("Starting a long task... please wait.")
    time.sleep(10) # Läuft in einem separaten Thread
    ctx.reply("Long task completed!")

Threadsicherheit

Funktionen mit threaded=True dürfen nicht direkt mit dem Reticulum Network Stack (RNS) oder Komponenten interagieren, die auf lxmfy.transport.py aufbauen, da diese im Allgemeinen nicht threadsicher sind. Verwenden Sie ctx.reply(), um aus einem Threaded-Befehl heraus Nachrichten an den Benutzer zu senden.

Ereignisse

Ereignissystem zur Verarbeitung verschiedener Bot-Ereignisse:

@bot.events.on("message_received", EventPriority.HIGHEST)
def handle_message(event):
    # Nachrichtenereignis verarbeiten
    pass

Testen

Die Projekt-Tests umfassen Zuverlässigkeits- und Stresstest-Szenarien in der Testsuite des Repositorys. Führen Sie sie mit dem Test-Runner des Repositorys aus.

Erweiterte Zuverlässigkeits-Testsuite

Das Framework enthält eine umfangreiche Suite automatisierter Tests für raue Umgebungsbedingungen:

  • Manifold Testing: Validiert die mathematische Topologie des NLP-Intent-Vektorraums.
  • Chaos Engineering: Simuliert Bitrott, SD-Karten-Ausfall und Speicherkorruption.
  • Temporal Drift: Prüft die Robustheit gegenüber Systemuhrsprüngen (±1 Jahr).
  • Leak Detection: Langzeit-Überwachung von Speicher, Dateideskriptoren und Threads.

Berechtigungen

Berechtigungssystem zur Zugriffskontrolle auf Bot-Funktionen:

from lxmfy import DefaultPerms

@bot.command(name="admin", description="Admin command", admin_only=True)
def admin_command(ctx):
    if ctx.is_admin:
        ctx.reply("Admin command executed")

Middleware

Middleware-System zur Verarbeitung von Nachrichten und Ereignissen:

@bot.middleware.register(MiddlewareType.PRE_COMMAND)
def pre_command_middleware(ctx):
    # Vor der Befehlsausführung verarbeiten
    pass

Anhänge

Unterstützung für das Senden von Dateien, Bildern und Audio:

from lxmfy import Attachment, AttachmentType

attachment = Attachment(
    type=AttachmentType.IMAGE,
    name="image.jpg",
    data=image_data,
    format="jpg"
)
bot.send_with_attachment(destination, "Here's an image", attachment)

Icon-Darstellung (LXMF-Feld)

Sie können für Ihren Bot ein eigenes Icon setzen, das kompatible LXMF-Clients anzeigen können. Dabei wird LXMF.FIELD_ICON_APPEARANCE verwendet.

from lxmfy import IconAppearance, pack_icon_appearance_field
import LXMF # Erforderlich für LXMF.FIELD_ICON_APPEARANCE

# Icon-Darstellung definieren
icon_data = IconAppearance(
    icon_name="smart_toy",  # Name aus Material Symbols
    fg_color=b'\xFF\xFF\xFF',  # Weißer Vordergrund (3 Bytes)
    bg_color=b'\x4A\x90\xE2'   # Blauer Hintergrund (3 Bytes)
)

# In das LXMF-Feldformat packen
icon_lxmf_field = pack_icon_appearance_field(icon_data)

# Nachricht mit diesem Icon senden
bot.send(
    destination_hash_str,
    "Hello from your friendly bot!",
    title="Bot Message",
    lxmf_fields=icon_lxmf_field
)

# Kann auch mit anderen Feldern kombiniert werden, z. B. Anhängen:
# attachment_field = pack_attachment(some_attachment)
# combined_fields = {**icon_lxmf_field, **attachment_field}
# bot.send(destination, "Message with icon and attachment", lxmf_fields=combined_fields)

Scheduler

System zur Aufgabenplanung:

@bot.scheduler.schedule(name="daily_task", cron_expr="0 0 * * *")
def daily_task():
    # Läuft täglich um Mitternacht
    pass

Signaturen

LXMFy bietet Konfigurationsoptionen für die in LXMF eingebaute kryptografische Nachrichtensignierung und -prüfung:

from lxmfy import LXMFBot

bot = LXMFBot(
    name="SecureBot",
    signature_verification_enabled=True,  # Signaturprüfung aktivieren
    require_message_signatures=False      # Auf True setzen, um unsignierte Nachrichten abzulehnen
)

Verarbeitung von Signaturen

LXMF führt alle kryptografischen Signatur- und Prüfoperationen automatisch mit RNS-Identitäten aus. Der SignatureManager von LXMFy ist eine Konfigurationsschicht, die:

  • steuert, ob die Signaturprüfung erzwungen wird
  • die Richtlinie für unsignierte Nachrichten festlegt (annehmen oder ablehnen)
  • sich in das Berechtigungssystem integriert (z. B. Prüfung für vertrauenswürdige Benutzer umgehen)

Die eigentlichen kryptografischen Operationen führt LXMF/RNS aus, nicht LXMFy.

Landlock-LSM-Sandbox

Auf Linux-Kerneln mit Landlock-Unterstützung (5.13+) kann LXMFy den Dateisystemzugriff für den Bot-Prozess und für externe Skript-Cogs einschränken.

Sandbox des Bot-Prozesses

Wenn landlock_enabled=True (Standard) und der Bot nicht im test_mode läuft, ruft der Bot apply_landlock_sandbox() während der Initialisierung auf. Systemverzeichnisse sind nur lesbar. Bot-Speicher, Konfiguration, Cogs, Reticulum-Konfiguration und Temp-Pfade bleiben schreibbar.

bot = LXMFBot(
    name="SecureBot",
    landlock_enabled=True,
)

status = bot.get_landlock_status()
# status-Schlüssel: landlock_kernel_supported, landlock_requested,
# landlock_auto_enabled, landlock_disabled_by_env, landlock_active

Überschreibung per Umgebungsvariable

  • LXMFY_LANDLOCK=0: Landlock auch auf unterstützten Kerneln deaktivieren
  • LXMFY_LANDLOCK=1: Landlock unter Linux unabhängig von der Auto-Erkennung versuchen
  • nicht gesetzt: landlock_enabled und Kernel-Auto-Erkennung folgen

Sandbox für externe Cogs

Skript-Cogs verwenden external_cogs_sandbox_type. Im Modus auto wird Landlock bevorzugt, wenn verfügbar, da keine externen Werkzeuge nötig sind. Die vollständige Liste der Sandbox-Optionen steht in der Anleitung Bots erstellen.

Identitäts-Pinning

LXMFy unterstützt optionales Identitäts-Pinning, um Identitätsdiebstahl zu verhindern, falls eine Identität rotiert oder kompromittiert wird. Wenn aktiviert, "pinnt" der Bot eine LXMF-Adresse an den zuerst gesehenen öffentlichen Schlüssel.

bot = LXMFBot(
    identity_pinning_enabled=True
)

SignatureManager-Methoden

Der SignatureManager ist als bot.signature_manager verfügbar, wenn signature_verification_enabled=True:

  • should_verify_message(sender): Feststellen, ob eine Nachricht des gegebenen Absenders geprüft werden soll
  • handle_unsigned_message(sender, message_hash): Nachrichten ohne gültige Signatur gemäß Richtlinie behandeln

Wie LXMF-Signaturen funktionieren

LXMF signiert alle ausgehenden Nachrichten automatisch mit der RNS-Identität des Absenders während der pack()-Operation. Beim Empfang validiert LXMF die Signaturen und stellt bereit:

  • message.signature_validated: Boolescher Wert, ob die Signatur gültig ist
  • message.unverified_reason: Grundcode, falls die Validierung fehlschlug (z. B. SIGNATURE_INVALID, SOURCE_UNKNOWN)

LXMFy nutzt diese eingebauten LXMF-Eigenschaften, um die Signaturrichtlinie Ihres Bots durchzusetzen.

Nachrichtenzustellung

LXMFy bietet erweiterte Zustellfunktionen, darunter Propagationsknoten und automatische Wiederholungen:

Propagationsknoten

Nachrichten über bestimmte Propagationsknoten senden, um die Zustellzuverlässigkeit im Reticulum-Netzwerk zu erhöhen:

# Propagationsknoten einmalig auf Konfigurations-/Laufzeitebene setzen
bot.set_propagation_node("<propagation_node_hash>")

# Mit konfiguriertem Zustellverhalten senden
bot.send(
    destination_hash,
    "Message content"
)

# Der Propagationsknoten-Hash sollte ein gültiger LXMF-Propagationsknoten
# im Reticulum-Netzwerk sein

Automatische Wiederholungen

Automatische Wiederholungsversuche für fehlgeschlagene Direktzustellungen konfigurieren:

bot = LXMFBot(
    name="ReliableBot",
    direct_delivery_retries=5,  # Direktzustellung bis zu 5-mal wiederholen
    propagation_fallback_enabled=True
)

bot.send(destination_hash, "Important message")

# direct_delivery_retries ist standardmäßig 3
# Die Wiederholungslogik verarbeitet Zustell-Callbacks automatisch

Das Wiederholungssystem verfolgt Zustellversuche pro Ziel und wiederholt fehlgeschlagene Zustellungen automatisch. Erfolgreiche Zustellungen setzen den Zähler für dieses Ziel zurück.

Nachrichten-Persistenz

Ausgehende Nachrichten können auf die Festplatte persistiert werden, damit sie auch nach einem Neustart des Bots zugestellt werden. Persistenz ist standardmäßig aktiviert. Die Ausgangswarteschlange im Speicher ist begrenzt (message_queue_size, Standard 50) und verwirft bei Volllast die älteste Nachricht. Ungültige Ziel-Hashes werden nicht wiederhergestellt.

bot = LXMFBot(
    message_persistence_enabled=True,
    message_queue_size=50,
)

Nachrichten-Handler

LXMFy stellt Dekoratoren für verschiedene Arten eingehender Nachrichten bereit:

Handler für erste Nachrichten

Die erste Nachricht jedes Benutzers verarbeiten:

@bot.on_first_message()
def welcome_user(sender, message):
    content = message.content.decode("utf-8")
    bot.send(sender, f"Welcome! You said: {content}")
    return True  # True zurückgeben, um die weitere Verarbeitung zu stoppen

Allgemeiner Nachrichten-Handler

Alle eingehenden Nachrichten vor der Befehlsverarbeitung verarbeiten:

@bot.on_message()
def handle_all_messages(sender, message):
    content = message.content.decode("utf-8").strip()

    # Eigene Logik hier
    if content.startswith("echo:"):
        bot.send(sender, content[5:])
        return True  # Weitere Verarbeitung stoppen

    return False  # Weiter zur Befehlsverarbeitung

Nachrichten-Handler werden in dieser Reihenfolge aufgerufen: 1. Handler für erste Nachrichten (wenn dies die erste Nachricht dieses Absenders ist) 2. Allgemeine Nachrichten-Handler (registriert mit on_message()) 3. Befehlsverarbeitung (wenn die Nachricht mit dem Befehlspräfix beginnt)

Reticulum Relay Chat (RRC)

Bots können RRC-Hubs über RNS-Links mit CBOR-Envelopes beitreten. Paket: lxmfy.rrc.

BotConfig-Optionen

  • rrc_enabled (bool, Standard False): Konfigurierte Hubs beim Start verbinden
  • rrc_hubs (Liste von Hex-Hashes): Ziel-Hashes der Hubs
  • rrc_rooms (Liste von Strings): Räume, die nach WELCOME automatisch betreten werden
  • rrc_nick (String oder None): Nickname in HELLO und Raumnachrichten
  • rrc_dest_name (String, Standard "rrc.hub"): Zielname zum Aufbau der Hub-Destination
  • rrc_auto_reconnect (bool, Standard True): Nach Link-Verlust neu verbinden
  • rrc_persist_sessions (bool, Standard True): Hubs und Räume über Neustarts hinweg speichern
  • reticulum_config_dir (String oder None): Reticulum-Konfigurationsverzeichnis. Kann auch über LXMFY_RETICULUM_CONFIG_DIR gesetzt werden. Dieselbe Konfiguration wie MeshChatX verwenden (oft ~/.reticulum), damit Hub-Announces sichtbar sind.

Beispiel

from lxmfy import LXMFBot, RRCMessage

bot = LXMFBot(
    name="RoomBot",
    reticulum_config_dir="~/.reticulum",
    rrc_enabled=True,
    rrc_hubs=["664fc0e8d2e448658e37bb3f34e6c88f"],
    rrc_rooms=["general"],
    rrc_nick="RoomBot",
)

@bot.on_rrc
def on_rrc(event, client, payload):
    if event == "msg" and isinstance(payload, RRCMessage) and payload.mention:
        client.send_message(payload.room, f"Hi {payload.nick}")

# Laufzeit-API
# bot.connect_rrc(hub_hash, rooms=["general"])
# bot.rrc.send_message("general", "hello")
# bot.rrc.send_notice("general", "notice")
# bot.rrc.send_action("general", "waves")
# bot.rrc.join("ops")
# bot.rrc.part("ops")
# bot.rrc.status()
# bot.disconnect_rrc()

Exportierte Typen

  • RRCClient: Sitzung mit einem Hub
  • RRCManager: Manager für mehrere Hubs (bot.rrc)
  • RRCMessage: Payload eines Raumereignisses (kind, room, text, nick, src, mention, ...)
  • RRC_VERSION: Konstante der Wire-Protokollversion

Übliche Ereignisse, die an @bot.on_rrc-Handler übergeben werden: status, welcome, joined, parted, msg, notice, action, motd, error und rtt.

Vorlagen

Das Framework enthält mehrere einsatzbereite Bot-Vorlagen:

EchoBot

Einfacher Echo-Bot, der Nachrichten wiederholt:

from lxmfy.templates import EchoBot

bot = EchoBot()
bot.run()

NoteBot

Notizen-Bot mit JSON-Speicher:

from lxmfy.templates import NoteBot

bot = NoteBot()
bot.run()

ReminderBot

Erinnerungs-Bot mit SQLite-Speicher:

from lxmfy.templates import ReminderBot

bot = ReminderBot()
bot.run()

RRCBot

RRC-Raum-Bot, der konfigurierten Hubs beitritt und auf @mentions antwortet. Standardmäßig Hub 664fc0e8d2e448658e37bb3f34e6c88f, Raum #general und ~/.reticulum, falls verfügbar.

from lxmfy.templates import RRCBot

bot = RRCBot(
    hubs=["664fc0e8d2e448658e37bb3f34e6c88f"],
    rooms=["general"],
    nick="RRCBot",
    reticulum_config_dir="~/.reticulum",
)
bot.run()

CLI-Werkzeuge

Das Framework stellt Kommandozeilen-Werkzeuge für die Bot-Verwaltung bereit:

# Neuen Bot erstellen
lxmfy create mybot

# Bot aus Vorlage erstellen
lxmfy create --template echo mybot
lxmfy create --template rrc my_rrc_bot

# Vorlagen-Bot ausführen
lxmfy run echo
lxmfy run rrc

# Signaturprüfung mit einer Nachricht testen
lxmfy signatures test

# Signaturprüfung aktivieren
lxmfy signatures enable

# Signaturprüfung deaktivieren
lxmfy signatures disable

Fehlerbehandlung

Shutdown- und Laufzeitfehler rund um bot.run() abfangen:

try:
    bot.run()
except KeyboardInterrupt:
    bot.cleanup()
except Exception as e:
    logger.error(f"Error running bot: {str(e)}")

Modulreferenz

Wird aus den Docstrings des Quellcodes generiert.

Bases: AnnounceMixin, DispatchMixin, CogMixin, InboundMixin, OutboundMixin, LinkMixin, RRCMixin, PropagationMixin

Main bot class for handling LXMF messages and commands.

This class manages the bot's lifecycle, including: - Message routing and delivery - Command registration and execution - Cog (extension) loading and management - Spam protection - Admin privileges

Source code in lxmfy/core.py
 59
 60
 61
 62
 63
 64
 65
 66
 67
 68
 69
 70
 71
 72
 73
 74
 75
 76
 77
 78
 79
 80
 81
 82
 83
 84
 85
 86
 87
 88
 89
 90
 91
 92
 93
 94
 95
 96
 97
 98
 99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
217
218
219
220
221
222
223
224
225
226
227
228
229
230
231
232
233
234
235
236
237
238
239
240
241
242
243
244
245
246
247
248
249
250
251
252
253
254
255
256
257
258
259
260
261
262
263
264
265
266
267
268
269
270
271
272
273
274
275
276
277
278
279
280
281
282
283
284
285
286
287
288
289
290
291
292
293
294
295
296
297
298
299
300
301
302
303
304
305
306
307
308
309
310
311
312
313
314
315
316
317
318
319
320
321
322
323
324
325
326
327
328
329
330
331
332
333
334
335
336
337
338
339
340
341
342
343
344
345
346
347
348
349
350
351
352
353
354
355
356
357
358
359
360
361
362
363
364
365
366
367
368
369
370
371
372
373
374
375
376
377
378
379
380
381
382
383
384
385
386
387
388
389
390
391
392
393
394
395
396
397
398
399
400
401
402
403
404
405
406
407
408
409
410
411
412
413
414
415
416
417
418
419
420
421
422
423
424
425
426
427
428
429
430
431
432
433
434
435
436
437
438
439
440
441
442
443
444
445
446
447
448
449
450
451
452
453
454
455
456
457
458
459
460
461
462
463
464
465
466
467
468
469
470
471
472
473
474
475
476
477
478
479
480
481
482
483
484
485
486
487
488
489
490
491
492
493
494
495
496
497
498
499
500
501
502
503
504
505
506
507
508
509
510
511
512
513
514
515
516
517
518
519
520
521
522
523
524
525
526
527
528
529
class LXMFBot(
    AnnounceMixin,
    DispatchMixin,
    CogMixin,
    InboundMixin,
    OutboundMixin,
    LinkMixin,
    RRCMixin,
    PropagationMixin,
):
    """Main bot class for handling LXMF messages and commands.

    This class manages the bot's lifecycle, including:
    - Message routing and delivery
    - Command registration and execution
    - Cog (extension) loading and management
    - Spam protection
    - Admin privileges
    """

    def __init__(self, name: str | None = None, **kwargs: Any):
        """Initialize a new LXMFBot instance.

        Args:
            name: Optional bot name, same as the name config key.
            **kwargs: Override default configuration settings

        """
        if name is not None:
            kwargs["name"] = name
        self.config = BotConfig(**kwargs)
        self.commands = {}
        self.cogs = {}
        self.first_message_handlers = []
        self.message_handlers = []
        self.reaction_handlers = []
        self.delivery_callbacks = []
        self.receipts = []
        self._receive_lock = threading.Lock()
        self._delivery_lock = threading.Lock()
        self._last_pending_flush = 0.0
        queue_size = max(1, int(getattr(self.config, "message_queue_size", 50) or 50))
        self.queue = Queue(maxsize=queue_size)
        self.announce_time = 600
        self.router = None
        self.local = None
        self.logger = logging.getLogger("lxmfy")
        self._configure_logging()
        self.thread_pool = ThreadPoolExecutor(
            max_workers=5,
        )  # For offloading CPU-bound or blocking I/O tasks
        self.scheduler = TaskScheduler(self)  # Initialize the scheduler

        if self.config.config_path:
            self.config_path = self.config.config_path
        else:
            self.config_path = os.path.join(os.getcwd(), "config")

        os.makedirs(self.config_path, exist_ok=True)
        if self.config.test_mode and not self.config.reticulum_config_dir:
            self.reticulum_config_dir = os.path.abspath(self.config_path)
        else:
            self.reticulum_config_dir = resolve_reticulum_config_dir(
                self.config.reticulum_config_dir,
                self.config_path,
            )
        os.makedirs(self.reticulum_config_dir, exist_ok=True)
        if not self.config.test_mode and is_isolated_reticulum_dir(
            self.reticulum_config_dir,
            self.config_path,
        ):
            ensure_isolated_share_instance_disabled(self.reticulum_config_dir)

        if self.config.storage_type == "json":
            self.storage = Storage(JSONStorage(self.config.storage_path))
        elif self.config.storage_type == "sqlite":
            self.storage = Storage(SQLiteStorage(self.config.storage_path))
        elif self.config.storage_type == "memory":
            self.storage = Storage(MemoryStorage())
        else:
            raise ValueError(
                f"Unknown storage_type {self.config.storage_type!r}; "
                "expected 'json', 'sqlite', or 'memory'",
            )

        self.admins = set(self.config.admins or [])
        self.permissions = PermissionManager(
            storage=self.storage,
            enabled=self.config.permissions_enabled,
            admins=self.admins,
        )

        self.events = EventManager(self.storage)

        self._register_builtin_events()

        self.middleware = MiddlewareManager()

        self.cogs_dir = os.path.join(self.config_path, self.config.cogs_dir)
        if self.config.cogs_enabled:
            os.makedirs(self.cogs_dir, exist_ok=True)
            init_file = os.path.join(self.cogs_dir, "__init__.py")
            if not os.path.exists(init_file):
                open(init_file, "w", encoding="utf-8").close()

        self.transport = Transport(self, self.storage)
        self.delivery = DeliveryTracker(storage=self.storage)
        self.delivery.load_persisted()
        self.spam_protection = SpamProtection(
            storage=self.storage,
            bot=self,
            rate_limit=self.config.rate_limit,
            cooldown=self.config.cooldown,
            max_warnings=self.config.max_warnings,
            warning_timeout=self.config.warning_timeout,
        )

        self.delivery_attempts = {}
        self._load_delivery_attempts()
        self.conversations = ConversationManager(self)

        self.landlock_active = False
        if not self.config.test_mode:
            storage_dir = os.path.abspath(
                os.path.expanduser(self.config.storage_path),
            )
            self.landlock_active = apply_landlock_sandbox(
                storage_dir=storage_dir,
                reticulum_config_dir=self.reticulum_config_dir,
                config_dir=self.config_path,
                cogs_dir=self.cogs_dir,
                config_enabled=self.config.landlock_enabled,
            )

        identity_file = os.path.join(self.config_path, "identity")

        self._owns_reticulum = False
        if not self.config.test_mode:
            # Initialize Reticulum (will raise exception if already running)
            if RNS.Reticulum.get_instance() is None:
                try:
                    RNS.Reticulum(
                        configdir=self.reticulum_config_dir,
                        loglevel=self.config.loglevel,
                    )
                    self._owns_reticulum = True
                except OSError:
                    if RNS.Reticulum.get_instance() is None:
                        raise

            if not os.path.isfile(identity_file):
                RNS.log("No Primary Identity file found, creating new...", RNS.LOG_INFO)
                identity = RNS.Identity(True)
                identity.to_file(identity_file)
            identity = RNS.Identity.from_file(identity_file)
            if identity is None:
                raise RuntimeError(f"Failed to load identity from {identity_file}")
            self.identity = identity
            RNS.log("Loaded identity from file", RNS.LOG_INFO)

            self.router = LXMRouter(
                identity=self.identity,
                storagepath=self.config_path,
                autopeer=self.config.autopeer_propagation,
                autopeer_maxdepth=self.config.autopeer_maxdepth,
                enforce_stamps=self.config.require_stamps,
            )
            self.local = self.router.register_delivery_identity(
                self.identity,
                display_name=self.config.name,
                stamp_cost=self.config.stamp_cost,
            )
            if self.local is None:
                raise RuntimeError("Failed to register delivery identity")
            self._sync_delivery_display_name()
            self.router.register_delivery_callback(self._message_received)
            self.local.set_link_established_callback(
                self._delivery_link_established,
            )

            if self.config.pending_sends_enabled:
                RNS.Transport.register_announce_handler(
                    PendingSendAnnounceHandler(self),
                )

        self._configure_propagation()

        if self.local:
            RNS.log(
                f"LXMF Router ready to receive on: {RNS.prettyhexrep(self.local.hash)}",
                RNS.LOG_INFO,
            )
        else:
            # Test mode - create mock components
            if os.path.isfile(identity_file):
                identity = RNS.Identity.from_file(identity_file)
                if identity is None:
                    raise RuntimeError(
                        f"Failed to load identity from {identity_file}",
                    )
                self.identity = identity
            else:
                self.identity = RNS.Identity()  # Create a basic identity for testing
                if self.config.config_path:
                    self.identity.to_file(identity_file)

            self.router = None
            self.local = None

        self.announce_enabled = self.config.announce_enabled
        self.announce_time = self.config.announce

        if self.announce_enabled and not self.config.test_mode:
            if self.announce_time > 0:
                # Schedule the announce task. announce_now throttles on the
                # announce interval file, so a sub-minute cron step still
                # respects announce_time.
                minutes = max(1, self.announce_time // 60)
                self.scheduler.add_task(
                    "announce_task",
                    self.announce_now,
                    f"*/{minutes} * * * *",  # Convert seconds to minutes for cron
                )
            if self.config.announce_immediately:
                self.announce_now(force=True)
                RNS.log("Initial announce sent", RNS.LOG_INFO)

        self.hot_reloading = self.config.hot_reloading
        self.command_prefix = self.config.command_prefix

        self.help_system = HelpSystem(self)
        register_admin_commands(self)

        self.nlp = IntentClassifier(threshold=self.config.nlp_threshold)
        self.intents = {}  # {intent_name: callback}

        self.link_handlers = []
        self.links = {}  # {dest_hash: Link}

        self.rrc = None
        self.rrc_handlers = []
        self._init_rrc()

        self.signature_manager = SignatureManager(
            self,
            verification_enabled=self.config.signature_verification_enabled,
            require_signatures=self.config.require_message_signatures,
            request_unknown_identities=self.config.request_unknown_identities,
        )

        self._load_persisted_queue()

        if self.config.cogs_enabled:
            load_cogs_from_directory(self)

    def _configure_logging(self) -> None:
        level = self.config.log_level
        if level is None:
            return
        if isinstance(level, str):
            level = logging.getLevelNamesMapping().get(level.upper())
            if level is None:
                raise ValueError(
                    f"Unknown log_level {self.config.log_level!r}",
                )
        logger = logging.getLogger("lxmfy")
        if not logger.handlers:
            handler = logging.StreamHandler()
            handler.setFormatter(
                logging.Formatter(
                    "%(asctime)s [%(levelname)s] %(name)s: %(message)s",
                    "%H:%M:%S",
                ),
            )
            logger.addHandler(handler)
        logger.setLevel(level)

    def get_debugger(self):
        """Return a Debugger bound to this bot."""
        from .debugger import Debugger

        return Debugger(bot=self)

    def diagnose_destination(
        self,
        destination: str,
        *,
        request_path: bool = False,
        wait: float = 0.0,
    ) -> dict:
        """Probe path and identity for a destination hash.

        Args:
            destination: Hex destination hash.
            request_path: Whether to request a path if missing.
            wait: Seconds to wait for a path after requesting.

        Returns:
            Dict describing identity/path status and hints.

        """
        return (
            self.get_debugger()
            .probe_destination(
                destination,
                request_path=request_path,
                wait=wait,
            )
            .to_dict()
        )

    def diagnose_connectivity(
        self,
        destination: str | None = None,
        *,
        request_path: bool = False,
        wait: float = 0.0,
    ) -> dict:
        """Run a connectivity doctor report for this bot.

        Args:
            destination: Optional peer hash to include in send diagnosis.
            request_path: Request a path when probing destination.
            wait: Seconds to wait for path discovery.

        Returns:
            Structured doctor report dict.

        """
        return (
            self.get_debugger()
            .run_doctor(
                destination,
                request_path=request_path,
                wait=wait,
            )
            .to_dict()
        )

    def on_delivery_event(self, callback=None):
        """Subscribe to the outbound delivery event stream.

        Usable as bot.on_delivery_event(fn) or as a decorator
        @bot.on_delivery_event(). Subscribers receive event dicts with
        ts, stage (queued, deferred, dispatched, delivered, failed,
        cancelled, dropped), and destination/message_id/hash/reason
        fields when available.
        """
        if callback is None:
            return self.delivery.subscribe
        return self.delivery.subscribe(callback)

    def get_landlock_status(self) -> dict[str, bool]:
        """Return Landlock LSM sandbox availability and activation state."""
        return landlock_status_dict(
            active=self.landlock_active,
            config_enabled=self.config.landlock_enabled,
        )

    def run(self, delay=10):
        """Run the bot"""
        self.scheduler.start()  # Start the scheduler
        try:
            while True:
                # Process outgoing queue with a timeout to prevent hanging
                while not self.queue.empty():
                    try:
                        lxm = self.queue.get(block=False)
                    except Exception:
                        break
                    try:
                        if self.router:
                            self.router.handle_outbound(lxm)
                            self.delivery.record(
                                "dispatched",
                                destination=_delivery_dest_hex(lxm),
                                message_id=_delivery_id_hex(lxm),
                                hash=_delivery_hex(lxm),
                            )
                        self._persist_queue()
                    except Exception as e:
                        self.logger.exception(
                            "Outbound send failed, requeueing",
                        )
                        self.delivery.record(
                            "failed",
                            destination=_delivery_dest_hex(lxm),
                            reason=str(e)[:120],
                        )
                        if not self._enqueue_outbound(lxm):
                            self.logger.exception(
                                "Failed to requeue after outbound error"
                            )
                        break

                retry_interval = max(
                    0,
                    int(getattr(self.config, "pending_sends_retry", 300) or 0),
                )
                if retry_interval and (
                    time.time() - self._last_pending_flush >= retry_interval
                ):
                    self._last_pending_flush = time.time()
                    self._flush_pending_sends()

                time.sleep(delay)

        except KeyboardInterrupt:
            pass
        finally:
            self.cleanup()

    def request_page(
        self,
        destination_hash: str,
        page_path: str,
        field_data: dict | None = None,
    ) -> dict:
        """Request a page from a destination.

        Args:
            destination_hash: The destination hash.
            page_path: The path to the page.
            field_data: Optional field data to send with the request.

        Returns:
            The response from the destination.

        """
        try:
            dest_hash_bytes = bytes.fromhex(destination_hash)
            return self.transport.request_page(dest_hash_bytes, page_path, field_data)
        except Exception:
            self.logger.exception("Error requesting page")
            raise

    def cleanup(self):
        """Clean up resources."""
        RNS.log("Cleaning up LXMFBot...", RNS.LOG_DEBUG)
        try:
            self._persist_queue()
        except Exception:
            self.logger.exception("Failed to persist queue during cleanup")
        self.conversations.cancel_all()
        if hasattr(self, "rrc") and self.rrc:
            try:
                self.rrc.shutdown()
            except Exception:
                self.logger.exception("RRC shutdown failed")
        self.transport.cleanup()
        self.thread_pool.shutdown(wait=False)
        self.scheduler.stop()
        if hasattr(self, "router") and self.router:
            try:
                self.router.exit_handler()
            except Exception as e:
                self.logger.debug("Router exit handler failed: %s", e)

        # Ensure Reticulum exits cleanly, but only if this bot started it.
        # A shared instance may still be serving other bots or tests.
        if not self.config.test_mode and self._owns_reticulum:
            try:
                RNS.Reticulum.exit_handler()
            except Exception as e:
                self.logger.debug("Reticulum exit handler failed: %s", e)
        RNS.log("LXMFBot cleanup complete", RNS.LOG_DEBUG)

    def validate(self) -> str:
        """Run validation checks and return formatted results."""
        results = validate_bot(self)
        return format_validation_results(results)

__init__(name=None, **kwargs)

Initialize a new LXMFBot instance.

Parameters:

Name Type Description Default
name str | None

Optional bot name, same as the name config key.

None
**kwargs Any

Override default configuration settings

{}
Source code in lxmfy/core.py
def __init__(self, name: str | None = None, **kwargs: Any):
    """Initialize a new LXMFBot instance.

    Args:
        name: Optional bot name, same as the name config key.
        **kwargs: Override default configuration settings

    """
    if name is not None:
        kwargs["name"] = name
    self.config = BotConfig(**kwargs)
    self.commands = {}
    self.cogs = {}
    self.first_message_handlers = []
    self.message_handlers = []
    self.reaction_handlers = []
    self.delivery_callbacks = []
    self.receipts = []
    self._receive_lock = threading.Lock()
    self._delivery_lock = threading.Lock()
    self._last_pending_flush = 0.0
    queue_size = max(1, int(getattr(self.config, "message_queue_size", 50) or 50))
    self.queue = Queue(maxsize=queue_size)
    self.announce_time = 600
    self.router = None
    self.local = None
    self.logger = logging.getLogger("lxmfy")
    self._configure_logging()
    self.thread_pool = ThreadPoolExecutor(
        max_workers=5,
    )  # For offloading CPU-bound or blocking I/O tasks
    self.scheduler = TaskScheduler(self)  # Initialize the scheduler

    if self.config.config_path:
        self.config_path = self.config.config_path
    else:
        self.config_path = os.path.join(os.getcwd(), "config")

    os.makedirs(self.config_path, exist_ok=True)
    if self.config.test_mode and not self.config.reticulum_config_dir:
        self.reticulum_config_dir = os.path.abspath(self.config_path)
    else:
        self.reticulum_config_dir = resolve_reticulum_config_dir(
            self.config.reticulum_config_dir,
            self.config_path,
        )
    os.makedirs(self.reticulum_config_dir, exist_ok=True)
    if not self.config.test_mode and is_isolated_reticulum_dir(
        self.reticulum_config_dir,
        self.config_path,
    ):
        ensure_isolated_share_instance_disabled(self.reticulum_config_dir)

    if self.config.storage_type == "json":
        self.storage = Storage(JSONStorage(self.config.storage_path))
    elif self.config.storage_type == "sqlite":
        self.storage = Storage(SQLiteStorage(self.config.storage_path))
    elif self.config.storage_type == "memory":
        self.storage = Storage(MemoryStorage())
    else:
        raise ValueError(
            f"Unknown storage_type {self.config.storage_type!r}; "
            "expected 'json', 'sqlite', or 'memory'",
        )

    self.admins = set(self.config.admins or [])
    self.permissions = PermissionManager(
        storage=self.storage,
        enabled=self.config.permissions_enabled,
        admins=self.admins,
    )

    self.events = EventManager(self.storage)

    self._register_builtin_events()

    self.middleware = MiddlewareManager()

    self.cogs_dir = os.path.join(self.config_path, self.config.cogs_dir)
    if self.config.cogs_enabled:
        os.makedirs(self.cogs_dir, exist_ok=True)
        init_file = os.path.join(self.cogs_dir, "__init__.py")
        if not os.path.exists(init_file):
            open(init_file, "w", encoding="utf-8").close()

    self.transport = Transport(self, self.storage)
    self.delivery = DeliveryTracker(storage=self.storage)
    self.delivery.load_persisted()
    self.spam_protection = SpamProtection(
        storage=self.storage,
        bot=self,
        rate_limit=self.config.rate_limit,
        cooldown=self.config.cooldown,
        max_warnings=self.config.max_warnings,
        warning_timeout=self.config.warning_timeout,
    )

    self.delivery_attempts = {}
    self._load_delivery_attempts()
    self.conversations = ConversationManager(self)

    self.landlock_active = False
    if not self.config.test_mode:
        storage_dir = os.path.abspath(
            os.path.expanduser(self.config.storage_path),
        )
        self.landlock_active = apply_landlock_sandbox(
            storage_dir=storage_dir,
            reticulum_config_dir=self.reticulum_config_dir,
            config_dir=self.config_path,
            cogs_dir=self.cogs_dir,
            config_enabled=self.config.landlock_enabled,
        )

    identity_file = os.path.join(self.config_path, "identity")

    self._owns_reticulum = False
    if not self.config.test_mode:
        # Initialize Reticulum (will raise exception if already running)
        if RNS.Reticulum.get_instance() is None:
            try:
                RNS.Reticulum(
                    configdir=self.reticulum_config_dir,
                    loglevel=self.config.loglevel,
                )
                self._owns_reticulum = True
            except OSError:
                if RNS.Reticulum.get_instance() is None:
                    raise

        if not os.path.isfile(identity_file):
            RNS.log("No Primary Identity file found, creating new...", RNS.LOG_INFO)
            identity = RNS.Identity(True)
            identity.to_file(identity_file)
        identity = RNS.Identity.from_file(identity_file)
        if identity is None:
            raise RuntimeError(f"Failed to load identity from {identity_file}")
        self.identity = identity
        RNS.log("Loaded identity from file", RNS.LOG_INFO)

        self.router = LXMRouter(
            identity=self.identity,
            storagepath=self.config_path,
            autopeer=self.config.autopeer_propagation,
            autopeer_maxdepth=self.config.autopeer_maxdepth,
            enforce_stamps=self.config.require_stamps,
        )
        self.local = self.router.register_delivery_identity(
            self.identity,
            display_name=self.config.name,
            stamp_cost=self.config.stamp_cost,
        )
        if self.local is None:
            raise RuntimeError("Failed to register delivery identity")
        self._sync_delivery_display_name()
        self.router.register_delivery_callback(self._message_received)
        self.local.set_link_established_callback(
            self._delivery_link_established,
        )

        if self.config.pending_sends_enabled:
            RNS.Transport.register_announce_handler(
                PendingSendAnnounceHandler(self),
            )

    self._configure_propagation()

    if self.local:
        RNS.log(
            f"LXMF Router ready to receive on: {RNS.prettyhexrep(self.local.hash)}",
            RNS.LOG_INFO,
        )
    else:
        # Test mode - create mock components
        if os.path.isfile(identity_file):
            identity = RNS.Identity.from_file(identity_file)
            if identity is None:
                raise RuntimeError(
                    f"Failed to load identity from {identity_file}",
                )
            self.identity = identity
        else:
            self.identity = RNS.Identity()  # Create a basic identity for testing
            if self.config.config_path:
                self.identity.to_file(identity_file)

        self.router = None
        self.local = None

    self.announce_enabled = self.config.announce_enabled
    self.announce_time = self.config.announce

    if self.announce_enabled and not self.config.test_mode:
        if self.announce_time > 0:
            # Schedule the announce task. announce_now throttles on the
            # announce interval file, so a sub-minute cron step still
            # respects announce_time.
            minutes = max(1, self.announce_time // 60)
            self.scheduler.add_task(
                "announce_task",
                self.announce_now,
                f"*/{minutes} * * * *",  # Convert seconds to minutes for cron
            )
        if self.config.announce_immediately:
            self.announce_now(force=True)
            RNS.log("Initial announce sent", RNS.LOG_INFO)

    self.hot_reloading = self.config.hot_reloading
    self.command_prefix = self.config.command_prefix

    self.help_system = HelpSystem(self)
    register_admin_commands(self)

    self.nlp = IntentClassifier(threshold=self.config.nlp_threshold)
    self.intents = {}  # {intent_name: callback}

    self.link_handlers = []
    self.links = {}  # {dest_hash: Link}

    self.rrc = None
    self.rrc_handlers = []
    self._init_rrc()

    self.signature_manager = SignatureManager(
        self,
        verification_enabled=self.config.signature_verification_enabled,
        require_signatures=self.config.require_message_signatures,
        request_unknown_identities=self.config.request_unknown_identities,
    )

    self._load_persisted_queue()

    if self.config.cogs_enabled:
        load_cogs_from_directory(self)

cleanup()

Clean up resources.

Source code in lxmfy/core.py
def cleanup(self):
    """Clean up resources."""
    RNS.log("Cleaning up LXMFBot...", RNS.LOG_DEBUG)
    try:
        self._persist_queue()
    except Exception:
        self.logger.exception("Failed to persist queue during cleanup")
    self.conversations.cancel_all()
    if hasattr(self, "rrc") and self.rrc:
        try:
            self.rrc.shutdown()
        except Exception:
            self.logger.exception("RRC shutdown failed")
    self.transport.cleanup()
    self.thread_pool.shutdown(wait=False)
    self.scheduler.stop()
    if hasattr(self, "router") and self.router:
        try:
            self.router.exit_handler()
        except Exception as e:
            self.logger.debug("Router exit handler failed: %s", e)

    # Ensure Reticulum exits cleanly, but only if this bot started it.
    # A shared instance may still be serving other bots or tests.
    if not self.config.test_mode and self._owns_reticulum:
        try:
            RNS.Reticulum.exit_handler()
        except Exception as e:
            self.logger.debug("Reticulum exit handler failed: %s", e)
    RNS.log("LXMFBot cleanup complete", RNS.LOG_DEBUG)

diagnose_connectivity(destination=None, *, request_path=False, wait=0.0)

Run a connectivity doctor report for this bot.

Parameters:

Name Type Description Default
destination str | None

Optional peer hash to include in send diagnosis.

None
request_path bool

Request a path when probing destination.

False
wait float

Seconds to wait for path discovery.

0.0

Returns:

Type Description
dict

Structured doctor report dict.

Source code in lxmfy/core.py
def diagnose_connectivity(
    self,
    destination: str | None = None,
    *,
    request_path: bool = False,
    wait: float = 0.0,
) -> dict:
    """Run a connectivity doctor report for this bot.

    Args:
        destination: Optional peer hash to include in send diagnosis.
        request_path: Request a path when probing destination.
        wait: Seconds to wait for path discovery.

    Returns:
        Structured doctor report dict.

    """
    return (
        self.get_debugger()
        .run_doctor(
            destination,
            request_path=request_path,
            wait=wait,
        )
        .to_dict()
    )

diagnose_destination(destination, *, request_path=False, wait=0.0)

Probe path and identity for a destination hash.

Parameters:

Name Type Description Default
destination str

Hex destination hash.

required
request_path bool

Whether to request a path if missing.

False
wait float

Seconds to wait for a path after requesting.

0.0

Returns:

Type Description
dict

Dict describing identity/path status and hints.

Source code in lxmfy/core.py
def diagnose_destination(
    self,
    destination: str,
    *,
    request_path: bool = False,
    wait: float = 0.0,
) -> dict:
    """Probe path and identity for a destination hash.

    Args:
        destination: Hex destination hash.
        request_path: Whether to request a path if missing.
        wait: Seconds to wait for a path after requesting.

    Returns:
        Dict describing identity/path status and hints.

    """
    return (
        self.get_debugger()
        .probe_destination(
            destination,
            request_path=request_path,
            wait=wait,
        )
        .to_dict()
    )

get_debugger()

Return a Debugger bound to this bot.

Source code in lxmfy/core.py
def get_debugger(self):
    """Return a Debugger bound to this bot."""
    from .debugger import Debugger

    return Debugger(bot=self)

get_landlock_status()

Return Landlock LSM sandbox availability and activation state.

Source code in lxmfy/core.py
def get_landlock_status(self) -> dict[str, bool]:
    """Return Landlock LSM sandbox availability and activation state."""
    return landlock_status_dict(
        active=self.landlock_active,
        config_enabled=self.config.landlock_enabled,
    )

on_delivery_event(callback=None)

Subscribe to the outbound delivery event stream.

Usable as bot.on_delivery_event(fn) or as a decorator @bot.on_delivery_event(). Subscribers receive event dicts with ts, stage (queued, deferred, dispatched, delivered, failed, cancelled, dropped), and destination/message_id/hash/reason fields when available.

Source code in lxmfy/core.py
def on_delivery_event(self, callback=None):
    """Subscribe to the outbound delivery event stream.

    Usable as bot.on_delivery_event(fn) or as a decorator
    @bot.on_delivery_event(). Subscribers receive event dicts with
    ts, stage (queued, deferred, dispatched, delivered, failed,
    cancelled, dropped), and destination/message_id/hash/reason
    fields when available.
    """
    if callback is None:
        return self.delivery.subscribe
    return self.delivery.subscribe(callback)

request_page(destination_hash, page_path, field_data=None)

Request a page from a destination.

Parameters:

Name Type Description Default
destination_hash str

The destination hash.

required
page_path str

The path to the page.

required
field_data dict | None

Optional field data to send with the request.

None

Returns:

Type Description
dict

The response from the destination.

Source code in lxmfy/core.py
def request_page(
    self,
    destination_hash: str,
    page_path: str,
    field_data: dict | None = None,
) -> dict:
    """Request a page from a destination.

    Args:
        destination_hash: The destination hash.
        page_path: The path to the page.
        field_data: Optional field data to send with the request.

    Returns:
        The response from the destination.

    """
    try:
        dest_hash_bytes = bytes.fromhex(destination_hash)
        return self.transport.request_page(dest_hash_bytes, page_path, field_data)
    except Exception:
        self.logger.exception("Error requesting page")
        raise

run(delay=10)

Run the bot

Source code in lxmfy/core.py
def run(self, delay=10):
    """Run the bot"""
    self.scheduler.start()  # Start the scheduler
    try:
        while True:
            # Process outgoing queue with a timeout to prevent hanging
            while not self.queue.empty():
                try:
                    lxm = self.queue.get(block=False)
                except Exception:
                    break
                try:
                    if self.router:
                        self.router.handle_outbound(lxm)
                        self.delivery.record(
                            "dispatched",
                            destination=_delivery_dest_hex(lxm),
                            message_id=_delivery_id_hex(lxm),
                            hash=_delivery_hex(lxm),
                        )
                    self._persist_queue()
                except Exception as e:
                    self.logger.exception(
                        "Outbound send failed, requeueing",
                    )
                    self.delivery.record(
                        "failed",
                        destination=_delivery_dest_hex(lxm),
                        reason=str(e)[:120],
                    )
                    if not self._enqueue_outbound(lxm):
                        self.logger.exception(
                            "Failed to requeue after outbound error"
                        )
                    break

            retry_interval = max(
                0,
                int(getattr(self.config, "pending_sends_retry", 300) or 0),
            )
            if retry_interval and (
                time.time() - self._last_pending_flush >= retry_interval
            ):
                self._last_pending_flush = time.time()
                self._flush_pending_sends()

            time.sleep(delay)

    except KeyboardInterrupt:
        pass
    finally:
        self.cleanup()

validate()

Run validation checks and return formatted results.

Source code in lxmfy/core.py
def validate(self) -> str:
    """Run validation checks and return formatted results."""
    results = validate_bot(self)
    return format_validation_results(results)

Configuration settings for LXMFBot.

Attributes:

Name Type Description
name str

The name of the bot. Defaults to "LXMFBot".

announce int

The announce interval in seconds. Defaults to 600.

announce_immediately bool

Whether to announce immediately on startup. Defaults to True.

admins set

A set of admin identity hashes. Defaults to an empty set.

hot_reloading bool

Whether to enable hot reloading of cogs. Defaults to False.

rate_limit int

The maximum number of messages allowed per cooldown period. Defaults to 5.

cooldown int

The cooldown period in seconds. Defaults to 60.

max_warnings int

The maximum number of spam warnings before action is taken. Defaults to 3.

warning_timeout int

The duration in seconds for which a spam warning is active. Defaults to 300.

command_prefix str

The prefix for bot commands. Defaults to "/".

cogs_dir str

The directory to load cogs from. Defaults to "cogs".

cogs_enabled bool

Whether to enable cogs. Defaults to True.

permissions_enabled bool

Whether to enable the permission system. Defaults to False.

storage_type str

The type of storage to use ("json" or "sqlite"). Defaults to "json".

storage_path str

The path to the storage file or directory. Defaults to "data".

first_message_enabled bool

Whether to enable first message handling. Defaults to True.

event_logging_enabled bool

Whether to enable event logging. Defaults to True.

max_logged_events int

The maximum number of events to log. Defaults to 1000.

event_middleware_enabled bool

Whether to enable event middleware. Defaults to True.

announce_enabled bool

Whether to enable bot announcements. Defaults to True.

signature_verification_enabled bool

Whether to enable cryptographic signature verification for incoming messages. Defaults to False.

require_message_signatures bool

Whether to reject unsigned messages when signature verification is enabled. Defaults to False.

require_stamps bool

Whether to reject messages with invalid stamps. Defaults to False.

request_unknown_identities bool

Whether to request unknown identities from the network when a message is received from an unknown source. Defaults to False.

stamp_cost int

Stamp cost required for messages to this bot. Outbound cost still comes from the peer announce unless send() overrides it. None disables inbound stamps. Defaults to None.

include_tickets bool

Include reply tickets on outbound messages so peers that require stamps can answer without stamp generation. Defaults to True.

direct_delivery_retries int

Number of times to retry direct delivery before falling back to propagation. Defaults to 3.

propagation_fallback_enabled bool

Whether to use propagation nodes as fallback after direct delivery fails. Defaults to True.

propagation_node str

The destination hash of the outbound propagation node. If None and autopeer_propagation is True, automatically discovers nodes. Defaults to None.

autopeer_propagation bool

Whether to automatically discover and peer with propagation nodes from announces. Defaults to False.

autopeer_maxdepth int

Maximum hop depth for auto-peering with propagation nodes. None = no limit. Defaults to 4.

enable_propagation_node bool

Whether to run this bot as a propagation node. Defaults to False.

message_storage_limit_mb float

Maximum storage for propagation node messages in megabytes. Only applies when enable_propagation_node is True. Defaults to 500 MB.

config_path str

The path to the bot configuration directory. If None, defaults to "config" in the current working directory. Defaults to None.

reticulum_config_dir str

The Reticulum config directory used for RNS shared instance/auth state. If None, uses LXMFY_RETICULUM_CONFIG_DIR when set, otherwise discovers the user/system Reticulum config (/etc/reticulum, ~/.config/reticulum, ~/.reticulum), and only then falls back to config_path. Isolated bot configs force share_instance=No to avoid RPC digest rejection with NomadNet/Columba.

test_mode bool

Whether to run in test mode (skips RNS initialization). Defaults to False.

log_level str or int or None

Logging level for lxmfy's own logger ("DEBUG", "INFO", "WARNING", ...). None leaves logging untouched. Defaults to "INFO".

loglevel int

RNS log level (0-7) passed to RNS.Reticulum. None uses the Reticulum config file setting. Defaults to None.

pending_sends_enabled bool

Hold outbound messages when the destination identity is not yet known, retrying on announces and periodic sweeps. Defaults to True.

pending_sends_max int

Maximum deferred messages kept for unknown destinations. Oldest are dropped beyond this. Defaults to 200.

pending_sends_ttl int

Seconds a deferred message is kept before being dropped. Defaults to 604800 (7 days).

pending_sends_retry int

Minimum seconds between deferred-send sweeps in run(). Defaults to 300.

announce_display_name_file str

Optional filename under config_path whose UTF-8 contents override the bot display name for LXMF delivery announces. If unset, bot_display_name.txt is read when present. Otherwise name is used.

Source code in lxmfy/config.py
@dataclass
class BotConfig:
    """Configuration settings for LXMFBot.

    Attributes:
        name (str): The name of the bot. Defaults to "LXMFBot".
        announce (int): The announce interval in seconds. Defaults to 600.
        announce_immediately (bool): Whether to announce immediately on startup. Defaults to True.
        admins (set): A set of admin identity hashes. Defaults to an empty set.
        hot_reloading (bool): Whether to enable hot reloading of cogs. Defaults to False.
        rate_limit (int): The maximum number of messages allowed per cooldown period. Defaults to 5.
        cooldown (int): The cooldown period in seconds. Defaults to 60.
        max_warnings (int): The maximum number of spam warnings before action is taken. Defaults to 3.
        warning_timeout (int): The duration in seconds for which a spam warning is active. Defaults to 300.
        command_prefix (str): The prefix for bot commands. Defaults to "/".
        cogs_dir (str): The directory to load cogs from. Defaults to "cogs".
        cogs_enabled (bool): Whether to enable cogs. Defaults to True.
        permissions_enabled (bool): Whether to enable the permission system. Defaults to False.
        storage_type (str): The type of storage to use ("json" or "sqlite"). Defaults to "json".
        storage_path (str): The path to the storage file or directory. Defaults to "data".
        first_message_enabled (bool): Whether to enable first message handling. Defaults to True.
        event_logging_enabled (bool): Whether to enable event logging. Defaults to True.
        max_logged_events (int): The maximum number of events to log. Defaults to 1000.
        event_middleware_enabled (bool): Whether to enable event middleware. Defaults to True.
        announce_enabled (bool): Whether to enable bot announcements. Defaults to True.
        signature_verification_enabled (bool): Whether to enable cryptographic signature verification for incoming messages. Defaults to False.
        require_message_signatures (bool): Whether to reject unsigned messages when signature verification is enabled. Defaults to False.
        require_stamps (bool): Whether to reject messages with invalid stamps. Defaults to False.
        request_unknown_identities (bool): Whether to request unknown identities from the network when a message is received from an unknown source. Defaults to False.
        stamp_cost (int): Stamp cost required for messages to this bot.
            Outbound cost still comes from the peer announce unless send()
            overrides it. None disables inbound stamps. Defaults to None.
        include_tickets (bool): Include reply tickets on outbound messages so
            peers that require stamps can answer without stamp generation.
            Defaults to True.
        direct_delivery_retries (int): Number of times to retry direct delivery before falling back to propagation. Defaults to 3.
        propagation_fallback_enabled (bool): Whether to use propagation nodes as fallback after direct delivery fails. Defaults to True.
        propagation_node (str): The destination hash of the outbound propagation node. If None and autopeer_propagation is True, automatically discovers nodes. Defaults to None.
        autopeer_propagation (bool): Whether to automatically discover and peer with propagation nodes from announces. Defaults to False.
        autopeer_maxdepth (int): Maximum hop depth for auto-peering with propagation nodes. None = no limit. Defaults to 4.
        enable_propagation_node (bool): Whether to run this bot as a propagation node. Defaults to False.
        message_storage_limit_mb (float): Maximum storage for propagation node messages in megabytes. Only applies when enable_propagation_node is True. Defaults to 500 MB.
        config_path (str): The path to the bot configuration directory. If None, defaults to "config" in the current working directory. Defaults to None.
        reticulum_config_dir (str): The Reticulum config directory used for RNS shared instance/auth state. If None, uses LXMFY_RETICULUM_CONFIG_DIR when set, otherwise discovers the user/system Reticulum config (/etc/reticulum, ~/.config/reticulum, ~/.reticulum), and only then falls back to config_path. Isolated bot configs force share_instance=No to avoid RPC digest rejection with NomadNet/Columba.
        test_mode (bool): Whether to run in test mode (skips RNS initialization). Defaults to False.
        log_level (str or int or None): Logging level for lxmfy's own logger ("DEBUG", "INFO", "WARNING", ...). None leaves logging untouched. Defaults to "INFO".
        loglevel (int): RNS log level (0-7) passed to RNS.Reticulum. None uses the Reticulum config file setting. Defaults to None.
        pending_sends_enabled (bool): Hold outbound messages when the destination identity is not yet known, retrying on announces and periodic sweeps. Defaults to True.
        pending_sends_max (int): Maximum deferred messages kept for unknown destinations. Oldest are dropped beyond this. Defaults to 200.
        pending_sends_ttl (int): Seconds a deferred message is kept before being dropped. Defaults to 604800 (7 days).
        pending_sends_retry (int): Minimum seconds between deferred-send sweeps in run(). Defaults to 300.
        announce_display_name_file (str): Optional filename under config_path whose UTF-8 contents override the bot display name for LXMF delivery announces. If unset, ``bot_display_name.txt`` is read when present. Otherwise ``name`` is used.

    """

    name: str = "LXMFBot"
    announce: int = 600
    announce_immediately: bool = True
    admins: set[str] | None = None
    hot_reloading: bool = False
    rate_limit: int = 5
    cooldown: int = 60
    max_warnings: int = 3
    warning_timeout: int = 300
    command_prefix: str = "/"
    cogs_dir: str = "cogs"
    cogs_enabled: bool = True
    permissions_enabled: bool = False
    storage_type: str = "json"
    storage_path: str = "data"
    first_message_enabled: bool = True
    event_logging_enabled: bool = True
    max_logged_events: int = 1000
    event_middleware_enabled: bool = True
    announce_enabled: bool = True
    signature_verification_enabled: bool = False
    require_message_signatures: bool = False
    require_stamps: bool = False
    request_unknown_identities: bool = False
    stamp_cost: int | None = None
    include_tickets: bool = True
    direct_delivery_retries: int = 3
    propagation_fallback_enabled: bool = True
    propagation_node: str | None = None
    autopeer_propagation: bool = False
    autopeer_maxdepth: int = 4
    enable_propagation_node: bool = False
    message_storage_limit_mb: float = 500.0
    config_path: str | None = None
    reticulum_config_dir: str | None = None
    announce_display_name_file: str | None = None
    test_mode: bool = False
    log_level: str | int | None = "INFO"
    loglevel: int | None = None
    identity_pinning_enabled: bool = False
    message_persistence_enabled: bool = True
    message_queue_size: int = 50
    pending_sends_enabled: bool = True
    pending_sends_max: int = 200
    pending_sends_ttl: int = 604800
    pending_sends_retry: int = 300
    dynamic_cogs_enabled: bool = True
    external_cogs_enabled: bool = True
    external_cogs_sandbox_enabled: bool = True
    external_cogs_sandbox_type: str = (
        "auto"  # 'auto', 'landlock', 'bwrap', 'firejail', 'none'
    )
    external_cogs_timeout: int = 30
    landlock_enabled: bool = True
    nlp_enabled: bool = False
    nlp_threshold: float = 0.5
    link_support_enabled: bool = False
    opportunistic_sending: bool = True
    lxmf_commands_enabled: bool = True
    rrc_enabled: bool = False
    rrc_hubs: list[str] | None = None
    rrc_rooms: list[str] | None = None
    rrc_nick: str | None = None
    rrc_dest_name: str = "rrc.hub"
    rrc_auto_reconnect: bool = True
    rrc_persist_sessions: bool = True

    def __post_init__(self):
        """Post-initialization to ensure admins is a set."""
        if self.admins is None:
            self.admins = set()
        if self.reticulum_config_dir is None:
            self.reticulum_config_dir = os.environ.get("LXMFY_RETICULUM_CONFIG_DIR")
        if self.rrc_hubs is None:
            self.rrc_hubs = []
        if self.rrc_rooms is None:
            self.rrc_rooms = []

    def __str__(self):
        """Return a string representation of the BotConfig object."""
        return f"BotConfig(name={self.name}, announce={self.announce}, announce_immediately={self.announce_immediately}, admins={self.admins}, hot_reloading={self.hot_reloading}, rate_limit={self.rate_limit}, cooldown={self.cooldown}, max_warnings={self.max_warnings}, warning_timeout={self.warning_timeout}, command_prefix={self.command_prefix}, cogs_dir={self.cogs_dir}, cogs_enabled={self.cogs_enabled}, permissions_enabled={self.permissions_enabled}, storage_type={self.storage_type}, storage_path={self.storage_path}, first_message_enabled={self.first_message_enabled}, event_logging_enabled={self.event_logging_enabled}, max_logged_events={self.max_logged_events}, event_middleware_enabled={self.event_middleware_enabled}, announce_enabled={self.announce_enabled}, signature_verification_enabled={self.signature_verification_enabled}, require_message_signatures={self.require_message_signatures}, require_stamps={self.require_stamps}, request_unknown_identities={self.request_unknown_identities}, stamp_cost={self.stamp_cost}, test_mode={self.test_mode}, identity_pinning_enabled={self.identity_pinning_enabled}, message_persistence_enabled={self.message_persistence_enabled}, dynamic_cogs_enabled={self.dynamic_cogs_enabled})"

__post_init__()

Post-initialization to ensure admins is a set.

Source code in lxmfy/config.py
def __post_init__(self):
    """Post-initialization to ensure admins is a set."""
    if self.admins is None:
        self.admins = set()
    if self.reticulum_config_dir is None:
        self.reticulum_config_dir = os.environ.get("LXMFY_RETICULUM_CONFIG_DIR")
    if self.rrc_hubs is None:
        self.rrc_hubs = []
    if self.rrc_rooms is None:
        self.rrc_rooms = []

__str__()

Return a string representation of the BotConfig object.

Source code in lxmfy/config.py
def __str__(self):
    """Return a string representation of the BotConfig object."""
    return f"BotConfig(name={self.name}, announce={self.announce}, announce_immediately={self.announce_immediately}, admins={self.admins}, hot_reloading={self.hot_reloading}, rate_limit={self.rate_limit}, cooldown={self.cooldown}, max_warnings={self.max_warnings}, warning_timeout={self.warning_timeout}, command_prefix={self.command_prefix}, cogs_dir={self.cogs_dir}, cogs_enabled={self.cogs_enabled}, permissions_enabled={self.permissions_enabled}, storage_type={self.storage_type}, storage_path={self.storage_path}, first_message_enabled={self.first_message_enabled}, event_logging_enabled={self.event_logging_enabled}, max_logged_events={self.max_logged_events}, event_middleware_enabled={self.event_middleware_enabled}, announce_enabled={self.announce_enabled}, signature_verification_enabled={self.signature_verification_enabled}, require_message_signatures={self.require_message_signatures}, require_stamps={self.require_stamps}, request_unknown_identities={self.request_unknown_identities}, stamp_cost={self.stamp_cost}, test_mode={self.test_mode}, identity_pinning_enabled={self.identity_pinning_enabled}, message_persistence_enabled={self.message_persistence_enabled}, dynamic_cogs_enabled={self.dynamic_cogs_enabled})"

Represents a generic attachment.

Attributes:

Name Type Description
type AttachmentType

The type of the attachment (AttachmentType).

name str

The name of the attachment.

data bytes

The binary data of the attachment.

format str | None

Optional format specifier (e.g., "png" for images).

Source code in lxmfy/attachments.py
@dataclass
class Attachment:
    """Represents a generic attachment.

    Attributes:
        type: The type of the attachment (AttachmentType).
        name: The name of the attachment.
        data: The binary data of the attachment.
        format: Optional format specifier (e.g., "png" for images).

    """

    type: AttachmentType
    name: str
    data: bytes
    format: str | None = None

Bases: IntEnum

Enumerates the different types of attachments supported.

FILE: Represents a generic file attachment. IMAGE: Represents an image attachment. AUDIO: Represents an audio attachment.

Source code in lxmfy/attachments.py
class AttachmentType(IntEnum):
    """Enumerates the different types of attachments supported.

    FILE: Represents a generic file attachment.
    IMAGE: Represents an image attachment.
    AUDIO: Represents an audio attachment.
    """

    FILE = 0x05
    IMAGE = 0x06
    AUDIO = 0x07

Bases: Flag

Default permission set

Source code in lxmfy/permissions.py
class DefaultPerms(Flag):
    """Default permission set"""

    NONE = 0
    # Basic permissions
    USE_BOT = auto()
    SEND_MESSAGES = auto()
    USE_COMMANDS = auto()

    # Elevated permissions
    MANAGE_MESSAGES = auto()
    MANAGE_COMMANDS = auto()
    MANAGE_USERS = auto()

    # Special permissions
    BYPASS_RATELIMIT = auto()
    BYPASS_SPAM = auto()
    VIEW_ADMIN_COMMANDS = auto()

    # Event system permissions
    VIEW_EVENTS = auto()
    MANAGE_EVENTS = auto()
    BYPASS_EVENT_CHECKS = auto()

    # Combined permissions
    ALL = (
        USE_BOT
        | SEND_MESSAGES
        | USE_COMMANDS
        | MANAGE_MESSAGES
        | MANAGE_COMMANDS
        | MANAGE_USERS
        | BYPASS_RATELIMIT
        | BYPASS_SPAM
        | VIEW_ADMIN_COMMANDS
        | VIEW_EVENTS
        | MANAGE_EVENTS
        | BYPASS_EVENT_CHECKS
    )

Manages permissions, roles, and user assignments

Source code in lxmfy/permissions.py
@dataclass
class PermissionManager:
    """Manages permissions, roles, and user assignments"""

    storage: Any
    enabled: bool = False
    admins: set[str] | None = None
    default_role: Role = field(
        default_factory=lambda: Role(
            "user",
            DefaultPerms.USE_BOT
            | DefaultPerms.SEND_MESSAGES
            | DefaultPerms.USE_COMMANDS,
        ),
    )
    admin_role: Role = field(
        default_factory=lambda: Role("admin", DefaultPerms.ALL, priority=100),
    )

    def __post_init__(self):
        self.roles: dict[str, Role] = {
            "user": self.default_role,
            "admin": self.admin_role,
        }
        self.user_roles: dict[str, set[str]] = {}
        self.load_data()

    def load_data(self):
        """Load permission data from storage"""
        stored_roles = self.storage.get("permissions:roles", {})
        stored_user_roles = self.storage.get("permissions:user_roles", {})

        # Convert stored roles back to Role objects
        for role_name, role_data in stored_roles.items():
            if role_name not in ["user", "admin"]:  # Don't override default/admin
                self.roles[role_name] = Role(
                    name=role_data["name"],
                    permissions=DefaultPerms(role_data["permissions"]),
                    priority=role_data["priority"],
                    description=role_data.get("description"),
                )

        self.user_roles = {
            user: set(roles) for user, roles in stored_user_roles.items()
        }

    def save_data(self):
        """Save permission data to storage"""
        # Convert roles to serializable format
        roles_data = {
            name: {
                "name": role.name,
                "permissions": role.permissions.value,
                "priority": role.priority,
                "description": role.description,
            }
            for name, role in self.roles.items()
        }

        self.storage.set("permissions:roles", roles_data)
        self.storage.set(
            "permissions:user_roles",
            {user: list(roles) for user, roles in self.user_roles.items()},
        )

    def create_role(
        self,
        name: str,
        permissions: DefaultPerms,
        priority: int = 0,
        description: str | None = None,
    ) -> Role:
        """Create a new role"""
        if name in self.roles:
            raise ValueError(f"Role {name} already exists")

        role = Role(name, permissions, priority, description)
        self.roles[name] = role
        self.save_data()
        return role

    def delete_role(self, name: str) -> bool:
        """Delete a role"""
        if name in ["user", "admin"]:
            raise ValueError("Cannot delete default or admin roles")

        if name in self.roles:
            del self.roles[name]
            # Remove role from all users
            for user_roles in self.user_roles.values():
                user_roles.discard(name)
            self.save_data()
            return True
        return False

    def assign_role(self, user: str, role_name: str):
        """Assign a role to a user"""
        if role_name not in self.roles:
            raise ValueError(f"Role {role_name} does not exist")

        if user not in self.user_roles:
            self.user_roles[user] = {self.default_role.name}

        self.user_roles[user].add(role_name)
        self.save_data()

    def remove_role(self, user: str, role_name: str):
        """Remove a role from a user"""
        if (
            user in self.user_roles
            and role_name in self.user_roles[user]
            and role_name != self.default_role.name
        ):
            self.user_roles[user].remove(role_name)
            self.save_data()

    def get_user_permissions(self, user: str) -> DefaultPerms:
        """Get combined permissions for a user"""
        if self.admins and user in self.admins:
            return self.admin_role.permissions

        if user not in self.user_roles:
            return self.default_role.permissions

        perms = DefaultPerms.NONE
        for role_name in self.user_roles[user]:
            if role_name in self.roles:
                perms |= self.roles[role_name].permissions

        return perms

    def has_permission(self, user: str, permission: DefaultPerms) -> bool:
        """Check if user has specific permission"""
        if not self.enabled:
            return True
        user_perms = self.get_user_permissions(user)
        return (user_perms & permission) == permission

assign_role(user, role_name)

Assign a role to a user

Source code in lxmfy/permissions.py
def assign_role(self, user: str, role_name: str):
    """Assign a role to a user"""
    if role_name not in self.roles:
        raise ValueError(f"Role {role_name} does not exist")

    if user not in self.user_roles:
        self.user_roles[user] = {self.default_role.name}

    self.user_roles[user].add(role_name)
    self.save_data()

create_role(name, permissions, priority=0, description=None)

Create a new role

Source code in lxmfy/permissions.py
def create_role(
    self,
    name: str,
    permissions: DefaultPerms,
    priority: int = 0,
    description: str | None = None,
) -> Role:
    """Create a new role"""
    if name in self.roles:
        raise ValueError(f"Role {name} already exists")

    role = Role(name, permissions, priority, description)
    self.roles[name] = role
    self.save_data()
    return role

delete_role(name)

Delete a role

Source code in lxmfy/permissions.py
def delete_role(self, name: str) -> bool:
    """Delete a role"""
    if name in ["user", "admin"]:
        raise ValueError("Cannot delete default or admin roles")

    if name in self.roles:
        del self.roles[name]
        # Remove role from all users
        for user_roles in self.user_roles.values():
            user_roles.discard(name)
        self.save_data()
        return True
    return False

get_user_permissions(user)

Get combined permissions for a user

Source code in lxmfy/permissions.py
def get_user_permissions(self, user: str) -> DefaultPerms:
    """Get combined permissions for a user"""
    if self.admins and user in self.admins:
        return self.admin_role.permissions

    if user not in self.user_roles:
        return self.default_role.permissions

    perms = DefaultPerms.NONE
    for role_name in self.user_roles[user]:
        if role_name in self.roles:
            perms |= self.roles[role_name].permissions

    return perms

has_permission(user, permission)

Check if user has specific permission

Source code in lxmfy/permissions.py
def has_permission(self, user: str, permission: DefaultPerms) -> bool:
    """Check if user has specific permission"""
    if not self.enabled:
        return True
    user_perms = self.get_user_permissions(user)
    return (user_perms & permission) == permission

load_data()

Load permission data from storage

Source code in lxmfy/permissions.py
def load_data(self):
    """Load permission data from storage"""
    stored_roles = self.storage.get("permissions:roles", {})
    stored_user_roles = self.storage.get("permissions:user_roles", {})

    # Convert stored roles back to Role objects
    for role_name, role_data in stored_roles.items():
        if role_name not in ["user", "admin"]:  # Don't override default/admin
            self.roles[role_name] = Role(
                name=role_data["name"],
                permissions=DefaultPerms(role_data["permissions"]),
                priority=role_data["priority"],
                description=role_data.get("description"),
            )

    self.user_roles = {
        user: set(roles) for user, roles in stored_user_roles.items()
    }

remove_role(user, role_name)

Remove a role from a user

Source code in lxmfy/permissions.py
def remove_role(self, user: str, role_name: str):
    """Remove a role from a user"""
    if (
        user in self.user_roles
        and role_name in self.user_roles[user]
        and role_name != self.default_role.name
    ):
        self.user_roles[user].remove(role_name)
        self.save_data()

save_data()

Save permission data to storage

Source code in lxmfy/permissions.py
def save_data(self):
    """Save permission data to storage"""
    # Convert roles to serializable format
    roles_data = {
        name: {
            "name": role.name,
            "permissions": role.permissions.value,
            "priority": role.priority,
            "description": role.description,
        }
        for name, role in self.roles.items()
    }

    self.storage.set("permissions:roles", roles_data)
    self.storage.set(
        "permissions:user_roles",
        {user: list(roles) for user, roles in self.user_roles.items()},
    )

Manages scheduled tasks and background processes.

Source code in lxmfy/scheduler.py
class TaskScheduler:
    """Manages scheduled tasks and background processes."""

    def __init__(self, bot: Any):
        """Initialize the TaskScheduler.

        Args:
            bot: The bot instance.

        """
        self.bot = bot
        self.tasks: dict[str, ScheduledTask] = {}
        self.background_tasks: list[Thread] = []
        self.stop_event = Event()
        self.logger = logging.getLogger(__name__)

    def schedule(self, name: str, cron_expr: str):
        """Decorator to schedule a task.

        Args:
            name (str): The name of the task.
            cron_expr (str): The cron expression for the task.

        """

        def decorator(func):
            """Adds the task to the scheduler."""
            self.add_task(name, func, cron_expr)
            return func

        return decorator

    def add_task(self, name: str, callback: Callable, cron_expr: str):
        """Add a scheduled task.

        Args:
            name (str): The name of the task.
            callback (Callable): The function to execute when the task runs.
            cron_expr (str): A cron-style expression defining when the task should run.

        """
        self.tasks[name] = ScheduledTask(name, callback, cron_expr)

    def remove_task(self, name: str):
        """Remove a scheduled task.

        Args:
            name (str): The name of the task to remove.

        """
        self.tasks.pop(name, None)

    def start(self):
        """Start the scheduler."""
        self.stop_event.clear()
        scheduler_thread = Thread(target=self._scheduler_loop, daemon=True)
        scheduler_thread.start()
        self.background_tasks.append(scheduler_thread)

    def stop(self):
        """Stop the scheduler."""
        self.stop_event.set()
        for task in self.background_tasks:
            task.join()
        self.background_tasks.clear()

    def _scheduler_loop(self):
        """Main scheduler loop.  Checks and runs tasks based on their cron expressions."""
        while not self.stop_event.is_set():
            current_time = datetime.now()  # noqa: DTZ005 - cron matches local wall time

            for task in list(self.tasks.values()):
                try:
                    if task.should_run(current_time):
                        run_sync(task.callback)
                        task.last_run = current_time
                except Exception:
                    self.logger.exception("Error running task %s", task.name)

            self.stop_event.wait(
                max(0, 60 - datetime.now().second)  # noqa: DTZ005
            )

__init__(bot)

Initialize the TaskScheduler.

Parameters:

Name Type Description Default
bot Any

The bot instance.

required
Source code in lxmfy/scheduler.py
def __init__(self, bot: Any):
    """Initialize the TaskScheduler.

    Args:
        bot: The bot instance.

    """
    self.bot = bot
    self.tasks: dict[str, ScheduledTask] = {}
    self.background_tasks: list[Thread] = []
    self.stop_event = Event()
    self.logger = logging.getLogger(__name__)

add_task(name, callback, cron_expr)

Add a scheduled task.

Parameters:

Name Type Description Default
name str

The name of the task.

required
callback Callable

The function to execute when the task runs.

required
cron_expr str

A cron-style expression defining when the task should run.

required
Source code in lxmfy/scheduler.py
def add_task(self, name: str, callback: Callable, cron_expr: str):
    """Add a scheduled task.

    Args:
        name (str): The name of the task.
        callback (Callable): The function to execute when the task runs.
        cron_expr (str): A cron-style expression defining when the task should run.

    """
    self.tasks[name] = ScheduledTask(name, callback, cron_expr)

remove_task(name)

Remove a scheduled task.

Parameters:

Name Type Description Default
name str

The name of the task to remove.

required
Source code in lxmfy/scheduler.py
def remove_task(self, name: str):
    """Remove a scheduled task.

    Args:
        name (str): The name of the task to remove.

    """
    self.tasks.pop(name, None)

schedule(name, cron_expr)

Decorator to schedule a task.

Parameters:

Name Type Description Default
name str

The name of the task.

required
cron_expr str

The cron expression for the task.

required
Source code in lxmfy/scheduler.py
def schedule(self, name: str, cron_expr: str):
    """Decorator to schedule a task.

    Args:
        name (str): The name of the task.
        cron_expr (str): The cron expression for the task.

    """

    def decorator(func):
        """Adds the task to the scheduler."""
        self.add_task(name, func, cron_expr)
        return func

    return decorator

start()

Start the scheduler.

Source code in lxmfy/scheduler.py
def start(self):
    """Start the scheduler."""
    self.stop_event.clear()
    scheduler_thread = Thread(target=self._scheduler_loop, daemon=True)
    scheduler_thread.start()
    self.background_tasks.append(scheduler_thread)

stop()

Stop the scheduler.

Source code in lxmfy/scheduler.py
def stop(self):
    """Stop the scheduler."""
    self.stop_event.set()
    for task in self.background_tasks:
        task.join()
    self.background_tasks.clear()

A scheduled task with cron-style timing.

Attributes:

Name Type Description
name str

The name of the task.

callback Callable

The function to execute when the task runs.

cron_expr str

A cron-style expression defining when the task should run (min hour day month weekday).

last_run Optional[datetime]

The last time the task was run.

enabled bool

Whether the task is currently enabled.

Source code in lxmfy/scheduler.py
@dataclass
class ScheduledTask:
    """A scheduled task with cron-style timing.

    Attributes:
        name (str): The name of the task.
        callback (Callable): The function to execute when the task runs.
        cron_expr (str): A cron-style expression defining when the task should run (min hour day month weekday).
        last_run (Optional[datetime]): The last time the task was run.
        enabled (bool): Whether the task is currently enabled.

    """

    name: str
    callback: Callable
    cron_expr: str
    last_run: datetime | None = None
    enabled: bool = True

    def should_run(self, current_time: datetime) -> bool:
        """Check if the task should run at the given time.

        Args:
            current_time (datetime): The current datetime.

        Returns:
            bool: True if the task should run, False otherwise.

        """
        if not self.enabled:
            return False

        if self.last_run and current_time - self.last_run < timedelta(minutes=1):
            return False

        return self._match_cron(current_time)

    def _match_cron(self, dt: datetime) -> bool:
        """Match the datetime against the cron expression.

        Args:
            dt (datetime): The datetime to match.

        Returns:
            bool: True if the datetime matches the cron expression, False otherwise.

        """
        parts = self.cron_expr.split()
        if len(parts) != 5:
            return False

        minute, hour, day, month, weekday = parts

        return (
            self._match_field(minute, dt.minute, 0, 59)
            and self._match_field(hour, dt.hour, 0, 23)
            and self._match_field(day, dt.day, 1, 31)
            and self._match_field(month, dt.month, 1, 12)
            and ScheduledTask._match_field(weekday, dt.weekday(), 0, 6)
        )

    @staticmethod
    def _match_field(pattern: str, value: int, min_val: int, max_val: int) -> bool:
        """Match a cron field pattern.

        Args:
            pattern (str): The cron field pattern to match.
            value (int): The value to check against the pattern.
            min_val (int): The minimum allowed value.
            max_val (int): The maximum allowed value.

        Returns:
            bool: True if the value matches the pattern, False otherwise.

        """
        if pattern == "*":
            return True

        parts = pattern.split(",")
        for part in parts:
            try:
                if "-" in part:
                    start, end = map(int, part.split("-"))
                    if min_val <= start <= value <= end <= max_val:
                        return True
                elif "/" in part:
                    step = int(part.split("/")[1])
                    if step > 0 and value % step == 0:
                        return True
                elif int(part) == value:
                    return True
            except ValueError:
                continue

        return False

should_run(current_time)

Check if the task should run at the given time.

Parameters:

Name Type Description Default
current_time datetime

The current datetime.

required

Returns:

Name Type Description
bool bool

True if the task should run, False otherwise.

Source code in lxmfy/scheduler.py
def should_run(self, current_time: datetime) -> bool:
    """Check if the task should run at the given time.

    Args:
        current_time (datetime): The current datetime.

    Returns:
        bool: True if the task should run, False otherwise.

    """
    if not self.enabled:
        return False

    if self.last_run and current_time - self.last_run < timedelta(minutes=1):
        return False

    return self._match_cron(current_time)

Facade for the underlying storage backend.

Source code in lxmfy/storage.py
class Storage:
    """Facade for the underlying storage backend."""

    def __init__(self, backend: StorageBackend):
        """Initialize a new Storage instance.

        Args:
            backend: The storage backend to use.

        """
        self.backend = backend

    def get(self, key: str, default: Any = None) -> Any:
        """Retrieve a value from storage.

        Args:
            key: The key to retrieve.
            default: The default value to return if the key is not found.

        Returns:
            The value associated with the key, or the default value if not found.

        """
        value = self.backend.get(key, default)
        return deserialize_value(value)

    def set(self, key: str, value: Any) -> None:
        """Store a value in storage.

        Args:
            key: The key to store the value under.
            value: The value to store.

        """
        serialized = serialize_value(value)
        self.backend.set(key, serialized)

    def delete(self, key: str) -> None:
        """Delete a value from storage.

        Args:
            key: The key to delete.

        """
        self.backend.delete(key)

    def exists(self, key: str) -> bool:
        """Check if a key exists in storage.

        Args:
            key: The key to check.

        Returns:
            True if the key exists, False otherwise.

        """
        return self.backend.exists(key)

    def scan(self, prefix: str) -> list:
        """Scan for keys with a given prefix.

        Args:
            prefix: The prefix to scan for.

        Returns:
            A list of keys that start with the prefix.

        """
        return self.backend.scan(prefix)

    def get_role_data(self, role_name: str) -> dict:
        """Helper method for permission system.

        Args:
            role_name: The name of the role.

        Returns:
            The role data.

        """
        return self.get(f"roles:{role_name}", {})

    def set_role_data(self, role_name: str, data: dict):
        """Helper method for permission system.

        Args:
            role_name: The name of the role.
            data: The role data.

        """
        self.set(f"roles:{role_name}", data)

    def get_user_roles(self, user_hash: str) -> list[str]:
        """Helper method for permission system.

        Args:
            user_hash: The hash of the user.

        Returns:
            The list of roles for the user.

        """
        return self.get(f"user_roles:{user_hash}", [])

    def set_user_roles(self, user_hash: str, roles: list[str]):
        """Helper method for permission system.

        Args:
            user_hash: The hash of the user.
            roles: The list of roles for the user.

        """
        self.set(f"user_roles:{user_hash}", roles)

__init__(backend)

Initialize a new Storage instance.

Parameters:

Name Type Description Default
backend StorageBackend

The storage backend to use.

required
Source code in lxmfy/storage.py
def __init__(self, backend: StorageBackend):
    """Initialize a new Storage instance.

    Args:
        backend: The storage backend to use.

    """
    self.backend = backend

delete(key)

Delete a value from storage.

Parameters:

Name Type Description Default
key str

The key to delete.

required
Source code in lxmfy/storage.py
def delete(self, key: str) -> None:
    """Delete a value from storage.

    Args:
        key: The key to delete.

    """
    self.backend.delete(key)

exists(key)

Check if a key exists in storage.

Parameters:

Name Type Description Default
key str

The key to check.

required

Returns:

Type Description
bool

True if the key exists, False otherwise.

Source code in lxmfy/storage.py
def exists(self, key: str) -> bool:
    """Check if a key exists in storage.

    Args:
        key: The key to check.

    Returns:
        True if the key exists, False otherwise.

    """
    return self.backend.exists(key)

get(key, default=None)

Retrieve a value from storage.

Parameters:

Name Type Description Default
key str

The key to retrieve.

required
default Any

The default value to return if the key is not found.

None

Returns:

Type Description
Any

The value associated with the key, or the default value if not found.

Source code in lxmfy/storage.py
def get(self, key: str, default: Any = None) -> Any:
    """Retrieve a value from storage.

    Args:
        key: The key to retrieve.
        default: The default value to return if the key is not found.

    Returns:
        The value associated with the key, or the default value if not found.

    """
    value = self.backend.get(key, default)
    return deserialize_value(value)

get_role_data(role_name)

Helper method for permission system.

Parameters:

Name Type Description Default
role_name str

The name of the role.

required

Returns:

Type Description
dict

The role data.

Source code in lxmfy/storage.py
def get_role_data(self, role_name: str) -> dict:
    """Helper method for permission system.

    Args:
        role_name: The name of the role.

    Returns:
        The role data.

    """
    return self.get(f"roles:{role_name}", {})

get_user_roles(user_hash)

Helper method for permission system.

Parameters:

Name Type Description Default
user_hash str

The hash of the user.

required

Returns:

Type Description
list[str]

The list of roles for the user.

Source code in lxmfy/storage.py
def get_user_roles(self, user_hash: str) -> list[str]:
    """Helper method for permission system.

    Args:
        user_hash: The hash of the user.

    Returns:
        The list of roles for the user.

    """
    return self.get(f"user_roles:{user_hash}", [])

scan(prefix)

Scan for keys with a given prefix.

Parameters:

Name Type Description Default
prefix str

The prefix to scan for.

required

Returns:

Type Description
list

A list of keys that start with the prefix.

Source code in lxmfy/storage.py
def scan(self, prefix: str) -> list:
    """Scan for keys with a given prefix.

    Args:
        prefix: The prefix to scan for.

    Returns:
        A list of keys that start with the prefix.

    """
    return self.backend.scan(prefix)

set(key, value)

Store a value in storage.

Parameters:

Name Type Description Default
key str

The key to store the value under.

required
value Any

The value to store.

required
Source code in lxmfy/storage.py
def set(self, key: str, value: Any) -> None:
    """Store a value in storage.

    Args:
        key: The key to store the value under.
        value: The value to store.

    """
    serialized = serialize_value(value)
    self.backend.set(key, serialized)

set_role_data(role_name, data)

Helper method for permission system.

Parameters:

Name Type Description Default
role_name str

The name of the role.

required
data dict

The role data.

required
Source code in lxmfy/storage.py
def set_role_data(self, role_name: str, data: dict):
    """Helper method for permission system.

    Args:
        role_name: The name of the role.
        data: The role data.

    """
    self.set(f"roles:{role_name}", data)

set_user_roles(user_hash, roles)

Helper method for permission system.

Parameters:

Name Type Description Default
user_hash str

The hash of the user.

required
roles list[str]

The list of roles for the user.

required
Source code in lxmfy/storage.py
def set_user_roles(self, user_hash: str, roles: list[str]):
    """Helper method for permission system.

    Args:
        user_hash: The hash of the user.
        roles: The list of roles for the user.

    """
    self.set(f"user_roles:{user_hash}", roles)

Bases: StorageBackend

JSON file-based storage backend.

Source code in lxmfy/storage.py
class JSONStorage(StorageBackend):
    """JSON file-based storage backend."""

    def __init__(self, directory: str):
        """Initialize a new JSONStorage instance.

        Args:
            directory: The directory to store the JSON files in.

        """
        self.directory = Path(directory)
        self.directory.mkdir(parents=True, exist_ok=True)
        self.cache: dict[str, Any] = {}
        self.logger = logging.getLogger(__name__)

    def get(self, key: str, default: Any = None) -> Any:
        """Retrieve a value from storage.

        Args:
            key: The key to retrieve.
            default: The default value to return if the key is not found.

        Returns:
            The value associated with the key, or the default value if not found.

        """
        if key in self.cache:
            return self.cache[key]

        file_path = self.directory / f"{key}.json"
        try:
            if file_path.exists():
                with open(file_path) as f:
                    data = json.load(f)
                    self.cache[key] = data
                    return data
        except Exception:
            self.logger.exception("Error reading %s", key)
        return default

    def set(self, key: str, value: Any) -> None:
        """Store a value in storage.

        Args:
            key: The key to store the value under.
            value: The value to store.

        """
        file_path = self.directory / f"{key}.json"
        try:
            with open(file_path, "w") as f:
                json.dump(value, f, indent=2)
            self.cache[key] = value
        except Exception:
            self.logger.exception("Error writing %s", key)
            raise

    def delete(self, key: str) -> None:
        """Delete a value from storage.

        Args:
            key: The key to delete.

        """
        file_path = self.directory / f"{key}.json"
        try:
            if file_path.exists():
                file_path.unlink()
            self.cache.pop(key, None)
        except Exception:
            self.logger.exception("Error deleting %s", key)
            raise

    def exists(self, key: str) -> bool:
        """Check if a key exists in storage.

        Args:
            key: The key to check.

        Returns:
            True if the key exists, False otherwise.

        """
        return (self.directory / f"{key}.json").exists()

    def scan(self, prefix: str) -> list:
        """Scan for keys with a given prefix.

        Args:
            prefix: The prefix to scan for.

        Returns:
            A list of keys that start with the prefix.

        """
        results = []
        try:
            for file in self.directory.glob(f"{prefix}*.json"):
                key = file.stem
                if key.startswith(prefix):
                    results.append(key)
        except Exception:
            self.logger.exception("Error scanning with prefix %s", prefix)
        return results

__init__(directory)

Initialize a new JSONStorage instance.

Parameters:

Name Type Description Default
directory str

The directory to store the JSON files in.

required
Source code in lxmfy/storage.py
def __init__(self, directory: str):
    """Initialize a new JSONStorage instance.

    Args:
        directory: The directory to store the JSON files in.

    """
    self.directory = Path(directory)
    self.directory.mkdir(parents=True, exist_ok=True)
    self.cache: dict[str, Any] = {}
    self.logger = logging.getLogger(__name__)

delete(key)

Delete a value from storage.

Parameters:

Name Type Description Default
key str

The key to delete.

required
Source code in lxmfy/storage.py
def delete(self, key: str) -> None:
    """Delete a value from storage.

    Args:
        key: The key to delete.

    """
    file_path = self.directory / f"{key}.json"
    try:
        if file_path.exists():
            file_path.unlink()
        self.cache.pop(key, None)
    except Exception:
        self.logger.exception("Error deleting %s", key)
        raise

exists(key)

Check if a key exists in storage.

Parameters:

Name Type Description Default
key str

The key to check.

required

Returns:

Type Description
bool

True if the key exists, False otherwise.

Source code in lxmfy/storage.py
def exists(self, key: str) -> bool:
    """Check if a key exists in storage.

    Args:
        key: The key to check.

    Returns:
        True if the key exists, False otherwise.

    """
    return (self.directory / f"{key}.json").exists()

get(key, default=None)

Retrieve a value from storage.

Parameters:

Name Type Description Default
key str

The key to retrieve.

required
default Any

The default value to return if the key is not found.

None

Returns:

Type Description
Any

The value associated with the key, or the default value if not found.

Source code in lxmfy/storage.py
def get(self, key: str, default: Any = None) -> Any:
    """Retrieve a value from storage.

    Args:
        key: The key to retrieve.
        default: The default value to return if the key is not found.

    Returns:
        The value associated with the key, or the default value if not found.

    """
    if key in self.cache:
        return self.cache[key]

    file_path = self.directory / f"{key}.json"
    try:
        if file_path.exists():
            with open(file_path) as f:
                data = json.load(f)
                self.cache[key] = data
                return data
    except Exception:
        self.logger.exception("Error reading %s", key)
    return default

scan(prefix)

Scan for keys with a given prefix.

Parameters:

Name Type Description Default
prefix str

The prefix to scan for.

required

Returns:

Type Description
list

A list of keys that start with the prefix.

Source code in lxmfy/storage.py
def scan(self, prefix: str) -> list:
    """Scan for keys with a given prefix.

    Args:
        prefix: The prefix to scan for.

    Returns:
        A list of keys that start with the prefix.

    """
    results = []
    try:
        for file in self.directory.glob(f"{prefix}*.json"):
            key = file.stem
            if key.startswith(prefix):
                results.append(key)
    except Exception:
        self.logger.exception("Error scanning with prefix %s", prefix)
    return results

set(key, value)

Store a value in storage.

Parameters:

Name Type Description Default
key str

The key to store the value under.

required
value Any

The value to store.

required
Source code in lxmfy/storage.py
def set(self, key: str, value: Any) -> None:
    """Store a value in storage.

    Args:
        key: The key to store the value under.
        value: The value to store.

    """
    file_path = self.directory / f"{key}.json"
    try:
        with open(file_path, "w") as f:
            json.dump(value, f, indent=2)
        self.cache[key] = value
    except Exception:
        self.logger.exception("Error writing %s", key)
        raise

Bases: StorageBackend

SQLite database storage backend.

Source code in lxmfy/storage.py
class SQLiteStorage(StorageBackend):
    """SQLite database storage backend."""

    def __init__(self, database_path: str):
        """Initialize a new SQLiteStorage instance.

        Args:
            database_path: The path to the SQLite database file.

        """
        self.database_path = database_path
        self.cache: dict[str, Any] = {}
        self.logger = logging.getLogger(__name__)
        self._ensure_db_dir()
        self._init_db()

    def _ensure_db_dir(self):
        """Ensure the database directory exists."""
        db_path = Path(self.database_path)
        db_dir = db_path.parent
        try:
            db_dir.mkdir(parents=True, exist_ok=True)
        except Exception:
            self.logger.exception("Failed to create database directory %s", db_dir)
            raise

    def _init_db(self):
        """Initialize the database table."""
        try:
            with sqlite3.connect(self.database_path) as conn:
                conn.execute("""
                    CREATE TABLE IF NOT EXISTS key_value (
                        key TEXT PRIMARY KEY,
                        value TEXT,
                        type TEXT,
                        created_at TIMESTAMP DEFAULT CURRENT_TIMESTAMP,
                        updated_at TIMESTAMP DEFAULT CURRENT_TIMESTAMP
                    )
                """)
                conn.execute("""
                    CREATE INDEX IF NOT EXISTS idx_key_prefix ON key_value(key)
                """)
        except sqlite3.OperationalError:
            self.logger.exception(
                "Failed to initialize database at %s", self.database_path
            )
            raise
        except Exception:
            self.logger.exception("Unexpected error initializing database")
            raise

    def get(self, key: str, default: Any = None) -> Any:
        """Retrieve a value from storage.

        Args:
            key: The key to retrieve.
            default: The default value to return if the key is not found.

        Returns:
            The value associated with the key, or the default value if not found.

        """
        if key in self.cache:
            return self.cache[key]

        try:
            with sqlite3.connect(self.database_path) as conn:
                cursor = conn.execute(
                    "SELECT value FROM key_value WHERE key = ?",
                    (key,),
                )
                row = cursor.fetchone()
                if row:
                    try:
                        value = json.loads(row[0])
                        self.cache[key] = value
                        return value
                    except json.JSONDecodeError:
                        return row[0]
        except Exception:
            self.logger.exception("Error reading %s", key)
        return default

    def set(self, key: str, value: Any) -> None:
        """Store a value in storage.

        Args:
            key: The key to store the value under.
            value: The value to store.

        """
        try:
            if isinstance(value, (dict, list)):
                serialized = json.dumps(value)
            else:
                serialized = str(value)

            with sqlite3.connect(self.database_path) as conn:
                conn.execute(
                    """
                    INSERT OR REPLACE INTO key_value (key, value, type, updated_at)
                    VALUES (?, ?, ?, CURRENT_TIMESTAMP)
                """,
                    (key, serialized, type(value).__name__),
                )
            self.cache[key] = value
        except Exception:
            self.logger.exception("Error writing %s", key)
            raise

    def delete(self, key: str) -> None:
        """Delete a value from storage.

        Args:
            key: The key to delete.

        """
        try:
            with sqlite3.connect(self.database_path) as conn:
                conn.execute("DELETE FROM key_value WHERE key = ?", (key,))
            self.cache.pop(key, None)
        except Exception:
            self.logger.exception("Error deleting %s", key)
            raise

    def exists(self, key: str) -> bool:
        """Check if a key exists in storage.

        Args:
            key: The key to check.

        Returns:
            True if the key exists, False otherwise.

        """
        try:
            with sqlite3.connect(self.database_path) as conn:
                cursor = conn.execute("SELECT 1 FROM key_value WHERE key = ?", (key,))
                return cursor.fetchone() is not None
        except Exception:
            self.logger.exception("Error checking existence of %s", key)
            return False

    def scan(self, prefix: str) -> list:
        """Scan for keys with a given prefix.

        Args:
            prefix: The prefix to scan for.

        Returns:
            A list of keys that start with the prefix.

        """
        try:
            with sqlite3.connect(self.database_path) as conn:
                cursor = conn.execute(
                    "SELECT key FROM key_value WHERE key LIKE ? ORDER BY key",
                    (f"{prefix}%",),
                )
                return [row[0] for row in cursor.fetchall()]
        except Exception:
            self.logger.exception("Error scanning with prefix %s", prefix)
            return []

__init__(database_path)

Initialize a new SQLiteStorage instance.

Parameters:

Name Type Description Default
database_path str

The path to the SQLite database file.

required
Source code in lxmfy/storage.py
def __init__(self, database_path: str):
    """Initialize a new SQLiteStorage instance.

    Args:
        database_path: The path to the SQLite database file.

    """
    self.database_path = database_path
    self.cache: dict[str, Any] = {}
    self.logger = logging.getLogger(__name__)
    self._ensure_db_dir()
    self._init_db()

delete(key)

Delete a value from storage.

Parameters:

Name Type Description Default
key str

The key to delete.

required
Source code in lxmfy/storage.py
def delete(self, key: str) -> None:
    """Delete a value from storage.

    Args:
        key: The key to delete.

    """
    try:
        with sqlite3.connect(self.database_path) as conn:
            conn.execute("DELETE FROM key_value WHERE key = ?", (key,))
        self.cache.pop(key, None)
    except Exception:
        self.logger.exception("Error deleting %s", key)
        raise

exists(key)

Check if a key exists in storage.

Parameters:

Name Type Description Default
key str

The key to check.

required

Returns:

Type Description
bool

True if the key exists, False otherwise.

Source code in lxmfy/storage.py
def exists(self, key: str) -> bool:
    """Check if a key exists in storage.

    Args:
        key: The key to check.

    Returns:
        True if the key exists, False otherwise.

    """
    try:
        with sqlite3.connect(self.database_path) as conn:
            cursor = conn.execute("SELECT 1 FROM key_value WHERE key = ?", (key,))
            return cursor.fetchone() is not None
    except Exception:
        self.logger.exception("Error checking existence of %s", key)
        return False

get(key, default=None)

Retrieve a value from storage.

Parameters:

Name Type Description Default
key str

The key to retrieve.

required
default Any

The default value to return if the key is not found.

None

Returns:

Type Description
Any

The value associated with the key, or the default value if not found.

Source code in lxmfy/storage.py
def get(self, key: str, default: Any = None) -> Any:
    """Retrieve a value from storage.

    Args:
        key: The key to retrieve.
        default: The default value to return if the key is not found.

    Returns:
        The value associated with the key, or the default value if not found.

    """
    if key in self.cache:
        return self.cache[key]

    try:
        with sqlite3.connect(self.database_path) as conn:
            cursor = conn.execute(
                "SELECT value FROM key_value WHERE key = ?",
                (key,),
            )
            row = cursor.fetchone()
            if row:
                try:
                    value = json.loads(row[0])
                    self.cache[key] = value
                    return value
                except json.JSONDecodeError:
                    return row[0]
    except Exception:
        self.logger.exception("Error reading %s", key)
    return default

scan(prefix)

Scan for keys with a given prefix.

Parameters:

Name Type Description Default
prefix str

The prefix to scan for.

required

Returns:

Type Description
list

A list of keys that start with the prefix.

Source code in lxmfy/storage.py
def scan(self, prefix: str) -> list:
    """Scan for keys with a given prefix.

    Args:
        prefix: The prefix to scan for.

    Returns:
        A list of keys that start with the prefix.

    """
    try:
        with sqlite3.connect(self.database_path) as conn:
            cursor = conn.execute(
                "SELECT key FROM key_value WHERE key LIKE ? ORDER BY key",
                (f"{prefix}%",),
            )
            return [row[0] for row in cursor.fetchall()]
    except Exception:
        self.logger.exception("Error scanning with prefix %s", prefix)
        return []

set(key, value)

Store a value in storage.

Parameters:

Name Type Description Default
key str

The key to store the value under.

required
value Any

The value to store.

required
Source code in lxmfy/storage.py
def set(self, key: str, value: Any) -> None:
    """Store a value in storage.

    Args:
        key: The key to store the value under.
        value: The value to store.

    """
    try:
        if isinstance(value, (dict, list)):
            serialized = json.dumps(value)
        else:
            serialized = str(value)

        with sqlite3.connect(self.database_path) as conn:
            conn.execute(
                """
                INSERT OR REPLACE INTO key_value (key, value, type, updated_at)
                VALUES (?, ?, ?, CURRENT_TIMESTAMP)
            """,
                (key, serialized, type(value).__name__),
            )
        self.cache[key] = value
    except Exception:
        self.logger.exception("Error writing %s", key)
        raise