event-forwarder.lua 9.3 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344
  1. --[[
  2. event-forwarder.lua - WebSocket-based event forwarder for OpenWrt
  3. Copyright (c) 2026 Luis Rosales - MIT License
  4. Size: ~30KB with dependencies
  5. ]]
  6. local socket = require("socket")
  7. local http = require("socket.http")
  8. local ltn12 = require("ltn12")
  9. local json = require("json")
  10. -- ============================================================================
  11. -- CONFIGURATION
  12. -- ============================================================================
  13. local config = {
  14. server_url = "wss://your-server.com/ws",
  15. server_token = "CHANGE_ME",
  16. router_id = "",
  17. reconnect_delay = 5,
  18. max_retries = 10,
  19. ping_interval = 30,
  20. -- File paths
  21. buffer_file = "/tmp/event_buffer",
  22. pid_file = "/var/run/event-forwarder.pid",
  23. -- Event sources
  24. events = {
  25. wifi_connect = true,
  26. wifi_disconnect = true,
  27. dhcp_lease = true,
  28. interface_up = true,
  29. interface_down = true,
  30. }
  31. }
  32. -- ============================================================================
  33. -- LOGGING
  34. -- ============================================================================
  35. local LOG_TAG = "event-forwarder"
  36. local function log(level, msg)
  37. io.popen(string.format('logger -t "%s" -p user.%s "%s"', LOG_TAG, level, msg:gsub('"', '\\"'))):close()
  38. end
  39. local function log_info(msg) log("info", msg) end
  40. local function log_err(msg) log("err", msg) end
  41. local function log_debug(msg)
  42. if os.getenv("DEBUG") then log("debug", msg) end
  43. end
  44. -- ============================================================================
  45. -- BUFFER (STORE-AND-FORWARD)
  46. -- ============================================================================
  47. local buffer = {}
  48. function buffer.load()
  49. local f = io.open(config.buffer_file, "r")
  50. if not f then return end
  51. for line in f:lines() do
  52. if line and line ~= "" then
  53. table.insert(buffer, line)
  54. end
  55. end
  56. f:close()
  57. log_info("Buffer loaded: " .. #buffer .. " events")
  58. end
  59. function buffer.save()
  60. local f = io.open(config.buffer_file, "w")
  61. if not f then return end
  62. for _, line in ipairs(buffer) do
  63. f:write(line .. "\n")
  64. end
  65. f:close()
  66. end
  67. function buffer.add(event_json)
  68. table.insert(buffer, event_json)
  69. -- Prevent unbounded growth
  70. while #buffer > 500 do
  71. table.remove(buffer, 1)
  72. end
  73. buffer.save()
  74. log_debug("Buffered event, total: " .. #buffer)
  75. end
  76. function buffer.flush(ws)
  77. if #buffer == 0 then return end
  78. log_info("Flushing " .. #buffer .. " buffered events...")
  79. local i = 1
  80. while i <= #buffer do
  81. local event = buffer[i]
  82. local sent = ws:send(event)
  83. if sent then
  84. table.remove(buffer, i)
  85. log_debug("Sent buffered event")
  86. else
  87. i = i + 1
  88. end
  89. end
  90. buffer.save()
  91. end
  92. function buffer.clear()
  93. buffer = {}
  94. os.execute("rm -f " .. config.buffer_file)
  95. end
  96. -- ============================================================================
  97. -- WEBSOCKET CLIENT (Simple implementation)
  98. -- ============================================================================
  99. local ws = {
  100. sock = nil,
  101. connected = false,
  102. }
  103. function ws.connect(url)
  104. local sock = require("socket").tcp()
  105. sock:settimeout(10)
  106. -- Parse URL
  107. local protocol, host, path = url:match("^(wss?)://([^/]+)(.*)")
  108. if not host then
  109. log_err("Invalid URL: " .. url)
  110. return nil
  111. end
  112. -- Connect
  113. local ok, err = sock:connect(host, 443)
  114. if not ok then
  115. return nil, "Connection failed: " .. tostring(err)
  116. end
  117. -- TLS handshake (simplified - use stunnel or openssl for real TLS)
  118. -- For now using raw socket - works with stunnel/Proxy
  119. ws.sock = sock
  120. ws.connected = true
  121. return ws
  122. end
  123. function ws.send(data)
  124. if not ws.connected then
  125. buffer.add(data) -- Buffer instead of dropping
  126. return false
  127. end
  128. local frame = string.format(
  129. "\x81%s%02x%s",
  130. string.char(0x80 + #data),
  131. #data,
  132. data
  133. )
  134. local ok, err = ws.sock:send(frame)
  135. if not ok then
  136. ws.connected = false
  137. buffer.add(data) -- Save to buffer on failure
  138. return false
  139. end
  140. return true
  141. end
  142. function ws.close()
  143. if ws.sock then
  144. ws.sock:close()
  145. ws.sock = nil
  146. end
  147. ws.connected = false
  148. end
  149. -- ============================================================================
  150. -- EVENT GENERATORS
  151. -- ============================================================================
  152. local function build_event(event_type, payload)
  153. return string.format([[{
  154. "router_id": "%s",
  155. "hostname": "%s",
  156. "event_type": "%s",
  157. "timestamp": "%s",
  158. "payload": %s
  159. }]],
  160. config.router_id,
  161. get_hostname() or "unknown",
  162. event_type,
  163. os.date("!%Y-%m-%dT%H:%M:%SZ"),
  164. json.encode(payload or {})
  165. )
  166. end
  167. local function get_hostname()
  168. local f = io.popen("hostname")
  169. if not f then return "unknown" end
  170. local name = f:read("*a"):gsub("%s+$", "")
  171. f:close()
  172. return name
  173. end
  174. -- ============================================================================
  175. -- EVENT LISTENERS
  176. -- ============================================================================
  177. -- DHCP Lease Events
  178. function listen_dhcp()
  179. local lease_file = "/var/lib/dnsmasq/dnsmasq.leases"
  180. local old_leases = {}
  181. while true do
  182. local f = io.open(lease_file, "r")
  183. if f then
  184. local leases = {}
  185. for line in f:lines() do
  186. local timestamp, mac, ip, hostname = line:match("(%d+) (%S+) (%S+) (%S+)")
  187. if mac then
  188. leases[mac] = { ip = ip, hostname = hostname, time = tonumber(timestamp) }
  189. -- New lease?
  190. if not old_leases[mac] then
  191. local event = build_event("dhcp_lease", {
  192. mac = mac,
  193. ip = ip,
  194. hostname = hostname,
  195. action = "new"
  196. })
  197. log_info("New DHCP: " .. mac .. " -> " .. ip)
  198. buffer.add(event)
  199. end
  200. end
  201. end
  202. old_leases = leases
  203. f:close()
  204. end
  205. socket.sleep(5)
  206. end
  207. end
  208. -- Interface Events (via ubus)
  209. function listen_interfaces()
  210. local proc = io.popen("ubus -m listen network.interface 2>/dev/null")
  211. if not proc then return end
  212. for line in proc:lines() do
  213. if line then
  214. local event_type = line:match('"action":"([^"]+)')
  215. local device = line:match('"interface":"([^"]+)')
  216. if event_type and device then
  217. local event = build_event("interface_" .. event_type, {
  218. device = device,
  219. action = event_type
  220. })
  221. log_info("Interface " .. event_type .. ": " .. device)
  222. buffer.add(event)
  223. end
  224. end
  225. end
  226. proc:close()
  227. end
  228. -- ============================================================================
  229. -- MAIN LOOP WITH RECONNECT
  230. -- ============================================================================
  231. local function main()
  232. -- Load config
  233. local uci_cursor = require("luci.model.uci").cursor()
  234. config.server_url = uci_cursor:get("event-forwarder", "server", "url") or config.server_url
  235. config.server_token = uci_cursor:get("event-forwarder", "server", "token") or config.server_token
  236. config.router_id = uci_cursor:get("event-forwarder", "router", "id") or get_hostname()
  237. log_info("Starting event forwarder...")
  238. log_info("Router ID: " .. config.router_id)
  239. log_info("Server: " .. config.server_url)
  240. -- Load buffered events
  241. buffer.load()
  242. -- Write PID
  243. local pidf = io.open(config.pid_file, "w")
  244. if pidf then
  245. pidf:write(tostring(os.getpid()))
  246. pidf:close()
  247. end
  248. -- Initial connect
  249. local retries = 0
  250. while true do
  251. -- Try to connect
  252. log_info("Connecting to server...")
  253. local ws, err = ws.connect(config.server_url)
  254. if ws and ws.connected then
  255. log_info("Connected!")
  256. retries = 0
  257. -- Flush buffer on connect
  258. buffer.flush(ws)
  259. -- Main event loop - in real impl, would use select() for both sockets
  260. -- For simplicity, just send buffered events periodically
  261. while ws.connected do
  262. socket.sleep(config.ping_interval)
  263. -- Send ping / keepalive
  264. if ws.connected then
  265. buffer.flush(ws)
  266. end
  267. end
  268. else
  269. log_err("Connection failed: " .. tostring(err))
  270. retries = retries + 1
  271. if retries >= config.max_retries then
  272. log_err("Max retries reached, resetting")
  273. retries = 0
  274. end
  275. end
  276. -- Wait before reconnect
  277. ws.close()
  278. log_info("Reconnecting in " .. config.reconnect_delay .. "s...")
  279. socket.sleep(config.reconnect_delay)
  280. end
  281. end
  282. -- Run
  283. main()