-- factoryos.net.node — FNP/1 node runtime.
--
-- Owns: identity, transport (modem), route table, service registry,
-- request/response with correlation + timeouts, announce/heartbeat,
-- per-message dispatch. Runs on the sched scheduler — node:run() blocks.
--
--   local node = Node.new({ name="pwr-01", role="controller", ... })
--   node:open()
--   node:service("power.read", { ops = { query = fn } })
--   node:spawn(function() ... end)   -- app tasks
--   node:run()

local sched    = require("factoryos.core.sched")
local envelope = require("factoryos.net.envelope")
local channels = require("factoryos.net.channels")
local util     = require("factoryos.util")
local log      = require("factoryos.core.log")

local Node = {}
Node.__index = Node

local BROADCAST = 65535

function Node.new(cfg)
  local self = setmetatable({}, Node)
  self.cfg       = cfg or {}
  self.name      = self.cfg.name or ("node-" .. tostring(os.getComputerID and os.getComputerID() or "?"))
  self.role      = self.cfg.role or "node"
  self.facility  = self.cfg.facility or "default"
  self.chan      = util.defaults(self.cfg.channels, {
    backbone = channels.backbone, telemetry = channels.telemetry,
    ui = channels.ui, display = channels.display,
  })
  self.cid       = self.cfg.cid or (os.getComputerID and os.getComputerID()) or 0
  self.hbInterval = self.cfg.hbInterval or 10
  self.idFile    = self.cfg.idFile or "/factoryos/node.id"
  self.token     = self.cfg.token
  self.tokenAdmin = self.cfg.token_admin

  self.id        = nil
  self.modem     = nil
  self.routes    = {}            -- nodeUUID -> reply channel (learned)
  self.peers     = {}            -- nodeUUID -> descriptor {id,name,role,facility,services,lastSeen}
  self.services  = {}            -- svc name -> {ops={op->{fn,level}|fn}, caps={}}
  self.handlers  = {}            -- event name -> {fns}
  self._pending  = {}            -- corr id -> coroutine
  self._timers   = {}            -- timer id -> corr id
  self._seen     = {}            -- dedup ring of message ids
  self._seenIdx  = 0
  return self
end

-- ── identity ──────────────────────────────────────────────────────────

function Node:_loadId()
  if self.cfg.id then self.id = self.cfg.id; return end
  if fs and fs.exists(self.idFile) then
    local h = fs.open(self.idFile, "r")
    if h then self.id = h.readAll(); h.close() end
    if self.id and #self.id > 0 then return end
  end
  self.id = util.uuid()
  if fs then
    local h = fs.open(self.idFile, "w")
    if h then h.write(self.id); h.close() end
  end
end

-- ── transport ─────────────────────────────────────────────────────────

--- Find and open the modem. Prefers wireless (ender/wireless); a cfg.modem
--- side/name overrides discovery. Opens backbone+telemetry+ui channels.
function Node:open()
  self:_loadId()
  local m
  if self.cfg.modem then
    m = peripheral.wrap(self.cfg.modem)
  else
    m = peripheral.find("modem", function(_, p)
      return p.isWireless and p.isWireless()
    end) or peripheral.find("modem")
  end
  if not m then return nil, "no modem peripheral" end
  self.modem = m
  self.modemSide = self.cfg.modem
    or (peripheral.getName and peripheral.getName(m)) or "modem"
  for _, ch in pairs(self.chan) do m.open(ch) end
  m.open(self.cid)     -- direct-address channel (routes resolve here)
  m.open(BROADCAST)
  return true
end

--- Transmit an envelope. `chan` selects the logical channel.
function Node:_send(env, chan)
  chan = chan or self.chan.backbone
  local dst = BROADCAST
  if env.to and env.to ~= "*" and self.routes[env.to] then
    dst = self.routes[env.to]
  end
  local s, err = envelope.encode(env)
  if not s then log.warn("encode failed: %s", err); return false, err end
  self.modem.transmit(dst, self.cid, s)
  return true
end

-- ── public messaging API ─────────────────────────────────────────────

function Node:service(name, def)
  self.services[name] = def
end

function Node:on(event, fn)
  self.handlers[event] = self.handlers[event] or {}
  table.insert(self.handlers[event], fn)
end

function Node:_emit_local(event, ...)
  for _, fn in ipairs(self.handlers[event] or {}) do pcall(fn, ...) end
end

function Node:descriptor()
  local svc = {}
  for name in pairs(self.services) do svc[#svc + 1] = name end
  table.sort(svc)
  return {
    id = self.id, name = self.name, role = self.role,
    facility = self.facility, services = svc,
    adopted = self.cfg.adopted == true,
  }
end

function Node:announce()
  self:_send(envelope.sys("announce", self.id, self:descriptor()))
end

function Node:hello()
  self:_send(envelope.sys("hello", self.id, self:descriptor()))
  self:announce()
end

--- Fire-and-forget event to a service. opts.chan overrides channel.
function Node:emit(svc, op, params, opts)
  opts = opts or {}
  local e = envelope.event(self.id, svc, op, params, { to = opts.to })
  return self:_send(e, opts.chan)
end

--- Request/response to a service (broadcast-addressed unless opts.to).
--- Runs in the calling task; blocks that task only. Returns ok, result.
function Node:call(svc, op, params, opts)
  opts = opts or {}
  -- local short-circuit: same-node services dispatch in-process
  local localSvc = self.services[svc]
  if localSvc and not opts.to then
    return self:_dispatch(localSvc, { svc = svc, op = op, p = params,
      auth = opts.auth, from = self.id, id = "local" })
  end
  local e = envelope.req(self.id, svc, op, params, { auth = opts.auth, to = opts.to })
  self._pending[e.id] = sched.current() or coroutine.running()
  local sent, err = self:_send(e, opts.chan)
  if not sent then
    self._pending[e.id] = nil
    return false, err
  end
  local tid = os.startTimer(opts.timeout or 5)
  self._timers[tid] = e.id
  local resp = sched.yield("__resp:" .. e.id)
  self._pending[e.id] = nil
  self._timers[tid] = nil
  if resp == nil then return false, "timeout" end
  local r = resp.p
  if type(r) ~= "table" then return false, "bad response" end
  return r.ok ~= false, r
end

-- ── dispatch ──────────────────────────────────────────────────────────

local LEVEL_OK = { read = 0, operate = 1, admin = 2 }

--- Unclaimed nodes (no node.cfg, or explicit adopted=false) answer admin
--- ops openly so a mainframe can adopt them. Once adopted=true is written,
--- normal token rules apply.
function Node:_claimable()
  return self.cfg._absent == true or self.cfg.adopted == false
end

function Node:_authLevel(env)
  if self:_claimable() then return 2 end
  if self.tokenAdmin and env.auth == self.tokenAdmin then return 2 end
  if self.token and env.auth == self.token then return 1 end
  return 0
end

--- Run a service op handler. Returns ok, resultTable.
function Node:_dispatch(svcDef, env)
  local spec = svcDef.ops and svcDef.ops[env.op]
  if not spec then return false, { ok = false, err = "no_such_op:" .. tostring(env.op) } end
  local fn, level = spec, "read"
  if type(spec) == "table" then fn, level = spec.fn, spec.level or "operate" end
  if LEVEL_OK[level] > self:_authLevel(env) then
    return false, { ok = false, err = "capability_denied" }
  end
  local ok, res = pcall(fn, env, env.p)
  if not ok then
    log.warn("svc %s op %s crashed: %s", env.svc, env.op, res)
    return false, { ok = false, err = "handler_error" }
  end
  if type(res) ~= "table" then res = { ok = true, result = res } end
  if res.ok == nil then res.ok = true end
  return true, res
end

function Node:_seenBefore(id)
  for _, v in ipairs(self._seen) do if v == id then return true end end
  return false
end

function Node:_markSeen(id)
  self._seenIdx = (self._seenIdx % 256) + 1
  self._seen[self._seenIdx] = id
end

function Node:_onMessage(replyChan, raw)
  local env = envelope.decode(raw)
  if not env then return end
  if env.from == self.id then return end           -- own broadcasts
  if self:_seenBefore(env.id) then return end
  self:_markSeen(env.id)

  self.routes[env.from] = replyChan
  self.lastRx = util.clock()
  local peer = self.peers[env.from] or {}
  peer.id, peer.lastSeen = env.from, util.clock()
  self.peers[env.from] = peer

  if env.to ~= "*" and env.to ~= self.id then return end

  if env.t == "announce" or env.t == "hello" then
    if type(env.p) == "table" then
      for k, v in pairs(env.p) do peer[k] = v end
      self:_emit_local("peer", peer)
    end
    if env.t == "hello" then self:announce() end

  elseif env.t == "resp" then
    local co = env.corr and self._pending[env.corr]
    if co then sched.wake(co, env) end

  elseif env.t == "req" or env.t == "event" then
    local svcDef = env.svc and self.services[env.svc]
    if svcDef then
      -- per-message task: a handler may itself call() remote services
      sched.spawn(function()
        local ok, res = self:_dispatch(svcDef, env)
        if env.t == "req" then
          local r = envelope.resp(self.id, env, res)
          if not ok and type(res) == "table" then r.p = res end
          self:_send(r)
        end
      end, "svc:" .. tostring(env.svc))
    else
      self:_emit_local("remote", env)  -- unknown service: observable anyway
      if env.t == "req" then
        local r = envelope.resp(self.id, env,
          { ok = false, err = "no_such_service:" .. tostring(env.svc) })
        self:_send(r)
      end
    end
  end
end

function Node:_onEvent(name, a, b, c, d)
  if name == "modem_message" then
    -- a=side b=channel c=replyChannel d=message e=distance
    if not self.modemSide or a == self.modemSide then
      self:_onMessage(c, d)
    end
  elseif name == "timer" then
    local corr = self._timers[a]
    local co = corr and self._pending[corr]
    if co then
      self._pending[corr] = nil
      sched.wake(co, nil)
    end
  else
    self:_emit_local(name, a, b, c, d)
  end
end

-- ── tasks & run ───────────────────────────────────────────────────────

function Node:spawn(fn, name, ...)
  return sched.spawn(fn, name or "node", ...)
end

function Node:_heartbeatLoop()
  while true do
    self:emit("core.registry", "heartbeat", self:descriptor())
    sched.sleep(self.hbInterval)
  end
end

--- Spawn dispatcher + heartbeat and announce. run() calls this; tests can
--- call it directly and pump sched.dispatch manually.
function Node:start()
  if self._started then return end
  self._started = true
  -- management agent on every node, before hello so it's advertised
  require("factoryos.services.agent").attach(self)
  self:spawn(function()
    while true do
      local name, a, b, c, d = sched.yield()
      self:_onEvent(name, a, b, c, d)
    end
  end, "dispatch")
  if not self.cfg.noHeartbeat then
    self:spawn(function() self:_heartbeatLoop() end, "heartbeat")
  end
  self:hello()
end

--- Start and run `main` (or idle) until terminate.
function Node:run(main)
  self:start()
  return sched.run(main or function()
    while true do sched.yield("__idle") end
  end)
end

return Node
