| 123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589590591592593594595596597598599600601602603604605606607608609610611612613614615616617618619620621622623624625626627628629630631632633634635636637638639640641642643644645646647648649650651652653654655656657658659660661662663664665666667668669670671672 |
- --[[
- client2server-luv - Lua WebSocket Event Forwarder for OpenWrt with luv async
-
- Features:
- - WebSocket connection to central server
- - Auto-reconnect on disconnect
- - Local buffer (store-and-forward while offline)
- - DHCP/WiFi/Interface event tracking
- - ASYNC monitoring via luv (libuv bindings)
- - Parallel polling for 50+ devices
-
- Copyright (c) 2026 Luis Rosales - MIT License
- ]]
- -- ============================================================================
- -- REQUIREMENTS
- -- ============================================================================
- -- Try to load luv (libuv bindings for async operations)
- local luv_ok, luv = pcall(require, "luv")
- -- Try to load websocket library, fallback to simple HTTP
- local ws_client = nil
- local has_websocket, websocket = pcall(require, "websocket")
- if has_websocket then
- ws_client = websocket.client.sync()
- end
- -- ============================================================================
- -- CONFIG
- -- ============================================================================
- local cfg = {
- url = os.getenv("SERVER_URL") or "wss://your-server.com/ws",
- token = os.getenv("SERVER_TOKEN") or "secret-token",
- router_id = os.getenv("ROUTER_ID") or "unknown",
- reconnect_delay = 5,
- ping_interval = 30,
- buffer_file = "/tmp/event_buffer",
- max_buffer = 100,
- poll_interval = 30, -- device poll interval (seconds)
- wan_poll_interval = 30, -- WAN poll interval (seconds)
- dhcp_poll_interval = 5, -- DHCP poll interval (seconds)
- poll_timeout = 5000, -- poll timeout (ms)
- }
- -- Load from UCI if available
- pcall(function()
- local uci = require("luci.model.uci").cursor()
- cfg.url = uci:get("event-forwarder", "server", "url") or cfg.url
- cfg.token = uci:get("event-forwarder", "server", "token") or cfg.token
- cfg.router_id = uci:get("event-forwarder", "router", "id") or cfg.router_id
- end)
- -- ============================================================================
- -- UTILITIES
- -- ============================================================================
- local function log(level, msg)
- os.execute(string.format('logger -t "client2server-luv" -p user.%s "%s" 2>/dev/null', level, msg))
- end
- local function get_hostname()
- -- Try multiple methods to get hostname
- local f = io.popen("cat /proc/sys/kernel/hostname 2>/dev/null || hostname 2>/dev/null || cat /etc/hostname 2>/dev/null || echo unknown")
- local h = f and f:read("*a") or "unknown"
- if f then f:close() end
- return (h:gsub("%s+$", ""))
- end
- cfg.router_id = cfg.router_id == "unknown" and get_hostname() or cfg.router_id
- local function json_encode(t)
- local parts = {}
- for k, v in pairs(t) do
- if type(v) == "string" then
- table.insert(parts, string.format('"%s": "%s"', k, v:gsub('"', '\\"')))
- elseif type(v) == "number" then
- table.insert(parts, string.format('"%s": %s', k, tostring(v)))
- elseif type(v) == "boolean" then
- table.insert(parts, string.format('"%s": %s', k, tostring(v)))
- end
- end
- return "{" .. table.concat(parts, ",") .. "}"
- end
- local function json_decode(str)
- local result = {}
- for k, v in str:gmatch('"([^"]+)":%s*([^},]+)') do
- v = v:gsub('[%s"]+', '')
- if v == "true" or v == "false" then
- result[k] = v == "true"
- elseif tonumber(v) then
- result[k] = tonumber(v)
- else
- result[k] = v
- end
- end
- return result
- end
- -- ============================================================================
- -- BUFFER (OFFLINE SUPPORT)
- -- ============================================================================
- local buffer = {
- events = {},
- dirty = false,
- }
- function buffer.load()
- local f = io.open(cfg.buffer_file, "r")
- if not f then return end
-
- for line in f:lines() do
- if line and line ~= "" then
- table.insert(buffer.events, line)
- end
- end
- f:close()
- log("info", "Loaded " .. #buffer.events .. " buffered events")
- end
- function buffer.save()
- local f = io.open(cfg.buffer_file, "w")
- if not f then return end
-
- for _, event in ipairs(buffer.events) do
- f:write(event .. "\n")
- end
- f:close()
- buffer.dirty = false
- end
- function buffer.add(json_event)
- table.insert(buffer.events, json_event)
-
- -- Limit buffer size
- while #buffer.events > cfg.max_buffer do
- table.remove(buffer.events, 1)
- end
-
- buffer.dirty = true
- end
- function buffer.flush(send_fn)
- if #buffer.events == 0 then return end
-
- local to_send = buffer.events
- buffer.events = {}
- buffer.dirty = false
-
- for _, event in ipairs(to_send) do
- local ok, err = pcall(send_fn, event)
- if not ok or err then
- -- Re-add to buffer on failure
- table.insert(buffer.events, event)
- end
- end
-
- buffer.save()
- end
- function buffer.clear()
- buffer.events = {}
- buffer.dirty = false
- os.remove(cfg.buffer_file)
- end
- -- ============================================================================
- -- WEBSOCKET (with fallback)
- -- ============================================================================
- local ws = {}
- function ws.connect(url)
- if ws_client then
- local sock, err = ws_client.connect(url)
- if err then
- log("err", "WS connect error: " .. err)
- return nil
- end
- return sock
- end
- return nil
- end
- function ws.send(sock, data)
- if sock then
- return sock:send(data)
- end
- return false, "no socket"
- end
- function ws.close(sock)
- if sock then
- sock:close()
- end
- end
- function ws.connected(sock)
- return sock ~= nil
- end
- function build_event(event_type, payload)
- return json_encode({
- type = event_type,
- router_id = cfg.router_id,
- timestamp = os.time(),
- payload = payload
- })
- end
- -- ============================================================================
- -- ASYNC POLLING WITH LUV
- -- ============================================================================
- -- Device polling results storage
- local poll_results = {}
- local poll_count = 0
- local poll_total = 0
- -- Poll a single device via ubus
- local function poll_single_device(device_id, device_ip)
- local cmd = string.format(
- 'ubus call network.interface.%s status 2>/dev/null',
- device_id
- )
-
- local f = io.popen(cmd)
- if not f then
- return { status = "error", ip = device_ip }
- end
-
- local result = f:read("*all")
- f:close()
-
- local parsed = json_decode(result)
- parsed.status = "online"
- parsed.ip = device_ip
-
- return parsed
- end
- -- Async parallel polling with luv
- local function poll_devices_parallel(device_list)
- if not luv then
- -- Fallback: sequential
- for _, dev in ipairs(device_list) do
- poll_results[dev.id] = poll_single_device(dev.id, dev.ip)
- end
- return
- end
-
- poll_results = {}
- poll_count = 0
- poll_total = #device_list
-
- log("info", "Starting parallel poll of " .. poll_total .. " devices")
-
- -- Poll each device in parallel using luv async
- for _, dev in ipairs(device_list) do
- local device_id = dev.id
- local device_ip = dev.ip
-
- -- Run in async task
- luv.new_task(function()
- local result = poll_single_device(device_id, device_ip)
- poll_results[device_id] = result
-
- poll_count = poll_count + 1
- if poll_count == poll_total then
- log("info", "All " .. poll_total .. " devices polled")
- -- Process results here if needed
- end
- end)
- end
- end
- -- ============================================================================
- -- WIFI EVENTS (hostapd via ubus)
- -- ============================================================================
- local function monitor_wifi_events()
- if not luv then
- log("warn", "luv not available for WiFi monitoring")
- return
- end
-
- -- Use luv to watch ubus for wireless events
- -- Note: ubus doesn't support event subscription directly,
- -- so we poll hostapd status periodically
-
- local timer = luv.new_timer()
- local last_clients = {}
-
- luv.timer_start(timer, 10000, 10000, function()
- -- Poll wireless clients
- local f = io.popen("ubus call hostapd.wlan0-1 get_clients 2>/dev/null")
- if f then
- local result = f:read("*all")
- f:close()
-
- if result and result ~= "" then
- -- Parse clients and detect changes
- local current_clients = {}
- for mac in result:gmatch('"([^"]+)":') do
- current_clients[mac] = true
- end
-
- -- Detect new connections
- for mac, _ in pairs(current_clients) do
- if not last_clients[mac] then
- local ev = build_event("wifi_connect", {
- mac = mac,
- interface = "wlan0-1"
- })
- log("info", "WiFi connected: " .. mac)
- buffer.add(ev)
- end
- end
-
- -- Detect disconnections
- for mac, _ in pairs(last_clients) do
- if not current_clients[mac] then
- local ev = build_event("wifi_disconnect", {
- mac = mac,
- interface = "wlan0-1"
- })
- log("info", "WiFi disconnected: " .. mac)
- buffer.add(ev)
- end
- end
-
- last_clients = current_clients
- end
- end
- end)
-
- log("info", "WiFi event monitor started")
- end
- -- ============================================================================
- -- DHCP LEASES (file watching with luv)
- -- ============================================================================
- local function monitor_dhcp_leases()
- if not luv then
- -- Fallback: sequential file polling
- monitor_dhcp_sequential()
- return
- end
-
- local lease_file = "/var/lib/dnsmasq/dnsmasq.leases"
- local old_leases = {}
- local last_mtime = 0
-
- -- Use luv fs_event to watch file changes
- local fs_event = luv.new_fs_event()
-
- luv.fs_event_start(fs_event, lease_file, function(err)
- if err then
- log("err", "DHCP fs_event error: " .. err)
- return
- end
-
- -- File changed, process leases
- process_leases(old_leases, function(mac, ip, hostname, action)
- local ev = build_event("dhcp_lease", {
- mac = mac,
- ip = ip,
- hostname = hostname,
- action = action
- })
- log("info", "DHCP: " .. action .. " - " .. mac .. " -> " .. ip)
- buffer.add(ev)
- end)
-
- old_leases = read_leases()
- end)
-
- log("info", "DHCP lease monitor started")
- end
- -- Helper: read current leases
- local function read_leases()
- local leases = {}
- local f = io.open("/var/lib/dnsmasq/dnsmasq.leases", "r")
- if not f then return leases end
-
- for line in f:lines() do
- local ts, mac, ip, name = line:match("(%d+)%s+(%S+)%s+(%S+)%s+(%S+)")
- if mac then
- leases[mac] = { ip = ip, hostname = name, time = tonumber(ts) }
- end
- end
- f:close()
- return leases
- end
- -- Helper: process leases and detect changes
- local function process_leases(old_leases, callback)
- local leases = read_leases()
-
- -- New leases
- for mac, info in pairs(leases) do
- if not old_leases[mac] then
- callback(mac, info.ip, info.hostname, "new")
- end
- end
-
- -- Expired leases
- for mac, info in pairs(old_leases) do
- if not leases[mac] then
- callback(mac, info.ip, info.hostname, "expired")
- end
- end
- end
- -- Fallback: sequential DHCP monitoring
- local function monitor_dhcp_sequential()
- local old_leases = {}
-
- while true do
- process_leases(old_leases, function(mac, ip, hostname, action)
- local ev = build_event("dhcp_lease", {
- mac = mac,
- ip = ip,
- hostname = hostname,
- action = action
- })
- log("info", "DHCP: " .. action .. " - " .. mac .. " -> " .. ip)
- buffer.add(ev)
- end)
-
- old_leases = read_leases()
- os.execute("sleep " .. cfg.dhcp_poll_interval)
- end
- end
- -- ============================================================================
- -- WAN MONITORING
- -- ============================================================================
- local last_wan_state = nil
- local function monitor_wan()
- if not luv then
- -- Fallback: sequential WAN polling
- monitor_wan_sequential()
- return
- end
-
- local timer = luv.new_timer()
-
- luv.timer_start(timer, cfg.wan_poll_interval * 1000, cfg.wan_poll_interval * 1000, function()
- local f = io.popen("ubus call network.interface.wan status 2>/dev/null")
- if f then
- local status = f:read("*all")
- f:close()
-
- local is_up = status:match('"up":%s*true') ~= nil
-
- if is_up ~= last_wan_state then
- local ev = build_event("wan_status", {
- up = is_up
- })
- log("info", "WAN: " .. (is_up and "up" or "down"))
- buffer.add(ev)
- last_wan_state = is_up
- end
- end
- end)
-
- log("info", "WAN monitor started")
- end
- -- Fallback: sequential WAN monitoring
- local function monitor_wan_sequential()
- while true do
- local f = io.popen("ubus call network.interface.wan status 2>/dev/null")
- if f then
- local status = f:read("*all")
- f:close()
-
- local is_up = status:match('"up":%s*true') ~= nil
-
- if is_up ~= last_wan_state then
- local ev = build_event("wan_status", {
- up = is_up
- })
- log("info", "WAN: " .. (is_up and "up" or "down"))
- buffer.add(ev)
- last_wan_state = is_up
- end
- end
-
- os.execute("sleep " .. cfg.wan_poll_interval)
- end
- end
- -- ============================================================================
- -- NETWORK STATUS POLLING
- -- ============================================================================
- local function poll_network_status()
- if not luv then
- -- Sequential fallback
- local f = io.popen("ubus call network getStatus 2>/dev/null")
- if f then f:close() end
- return
- end
-
- local timer = luv.new_timer()
-
- luv.timer_start(timer, cfg.poll_interval * 1000, cfg.poll_interval * 1000, function()
- local f = io.popen("ubus call network getStatus 2>/dev/null")
- if f then
- local status = f:read("*all")
- f:close()
-
- if status and status ~= "" then
- local ev = build_event("network_status", {
- status = status
- })
- buffer.add(ev)
- end
- end
- end)
-
- log("info", "Network status poller started")
- end
- -- ============================================================================
- -- MAIN EVENT LOOP
- -- ============================================================================
- local function main()
- log("info", "client2server-luv starting...")
- log("info", "Router: " .. cfg.router_id)
- log("info", "Server: " .. cfg.url)
-
- if luv_ok then
- log("info", "luv available - using async mode")
- else
- log("warn", "luv NOT available - using blocking mode")
- end
-
- -- Load buffered events
- buffer.load()
-
- -- Save PID
- local pf = io.open("/var/run/client2server-luv.pid", "w")
- if pf then
- local f = io.popen("echo $$")
- local pid = f and f:read("*a") or "0"
- if f then f:close() end
- pf:write(pid:gsub("%s+", ""))
- pf:close()
- end
-
- local sock = nil
- local retries = 0
-
- if luv_ok then
- -- Use luv event loop
- luv.run(function()
- -- Start all monitors in parallel
- monitor_wifi_events()
- monitor_dhcp_leases()
- monitor_wan()
- poll_network_status()
-
- -- Keep the event loop running
- local idle = luv.new_idle()
- luv.idle_start(idle, function()
- -- Idle work
- end)
-
- -- WebSocket connection and keepalive
- while true do
- log("info", "Connecting to server...")
-
- if ws_client then
- sock = ws.connect(cfg.url)
- end
-
- if sock then
- log("info", "Connected!")
- retries = 0
-
- -- Flush buffer
- buffer.flush(function(data)
- return ws.send(sock, data)
- end)
-
- -- Keep alive loop
- local loop_count = 0
- while ws.connected(sock) and loop_count < (cfg.ping_interval / 5) do
- luv.sleep(5000)
- loop_count = loop_count + 1
-
- -- Periodic flush
- buffer.flush(function(data)
- return ws.send(sock, data)
- end)
- end
- else
- log("err", "Connection failed")
- retries = retries + 1
- end
-
- -- Cleanup and reconnect
- ws.close(sock)
- sock = nil
-
- -- Save buffer on disconnect
- buffer.save()
-
- -- Delay before reconnect
- luv.sleep(cfg.reconnect_delay * 1000)
- end
- end)
- else
- -- Fallback: sequential mode
- while true do
- -- Sequential monitoring
- monitor_dhcp_sequential()
- monitor_wan_sequential()
-
- -- WebSocket connection
- log("info", "Connecting to server...")
-
- if ws_client then
- sock = ws.connect(cfg.url)
- end
-
- if sock then
- log("info", "Connected!")
- retries = 0
-
- buffer.flush(function(data)
- return ws.send(sock, data)
- end)
-
- local loop_count = 0
- while ws.connected(sock) and loop_count < (cfg.ping_interval / 5) do
- os.execute("sleep 5")
- loop_count = loop_count + 1
-
- buffer.flush(function(data)
- return ws.send(sock, data)
- end)
- end
- else
- log("err", "Connection failed")
- retries = retries + 1
- end
-
- ws.close(sock)
- buffer.save()
- os.execute("sleep " .. cfg.reconnect_delay)
- end
- end
- end
- -- ============================================================================
- -- START
- -- ============================================================================
- main()
|