Packages

ExClaw — OpenClaw rebuilt on ADK Elixir

Current section

Files

Jump to
ex_claw lib ex_claw workspace file_watcher.ex
Raw

lib/ex_claw/workspace/file_watcher.ex

defmodule ExClaw.Workspace.FileWatcher do
@moduledoc """
Watches workspace files for external modifications.
Uses polling (every 2 seconds by default) comparing file mtimes.
When a change is detected:
1. Validates the new content
2. If valid: updates cached state, broadcasts `{:file_changed, name, content}`
3. If invalid: logs warning, keeps old cached version, broadcasts `{:file_invalid, name, reason}`
No external dependencies — pure OTP.
"""
use GenServer
require Logger
alias ExClaw.Workspace.Validator
@default_poll_interval_ms 2_000
@file_atoms %{
"SOUL.md" => :soul,
"IDENTITY.md" => :identity,
"USER.md" => :user,
"AGENTS.md" => :agents,
"TOOLS.md" => :tools,
"HEARTBEAT.md" => :heartbeat,
"MEMORY.md" => :memory
}
# --- Public API ---
@doc "Start the file watcher GenServer."
@spec start_link(keyword()) :: GenServer.on_start()
def start_link(opts \\ []) do
GenServer.start_link(__MODULE__, opts, name: Keyword.get(opts, :name, __MODULE__))
end
@doc "Subscribe the calling process to file change notifications."
@spec subscribe(GenServer.server()) :: :ok
def subscribe(server \\ __MODULE__) do
GenServer.call(server, {:subscribe, self()})
end
@doc "Get cached content of a file by atom name."
@spec get(atom(), GenServer.server()) :: {:ok, String.t()} | {:error, :not_found}
def get(name, server \\ __MODULE__) do
GenServer.call(server, {:get, name})
end
@doc "Get all cached files as a map."
@spec get_all(GenServer.server()) :: map()
def get_all(server \\ __MODULE__) do
GenServer.call(server, :get_all)
end
# --- GenServer Callbacks ---
@impl GenServer
def init(opts) do
poll_interval = Keyword.get(opts, :poll_interval_ms, @default_poll_interval_ms)
root_dir = Keyword.get(opts, :root_dir) || workspace_root()
context_dir = Keyword.get(opts, :context_dir)
state = %{
root_dir: root_dir,
context_dir: context_dir,
poll_interval: poll_interval,
file_meta: %{},
cache: %{},
subscribers: MapSet.new()
}
# Initial load
state = load_all_files(state)
# Schedule first poll
schedule_poll(poll_interval)
{:ok, state}
end
@impl GenServer
def handle_call({:subscribe, pid}, _from, state) do
Process.monitor(pid)
{:reply, :ok, %{state | subscribers: MapSet.put(state.subscribers, pid)}}
end
def handle_call({:get, name}, _from, state) do
case Map.fetch(state.cache, name) do
{:ok, content} -> {:reply, {:ok, content}, state}
:error -> {:reply, {:error, :not_found}, state}
end
end
def handle_call(:get_all, _from, state) do
{:reply, state.cache, state}
end
@impl GenServer
def handle_info(:poll, state) do
state = poll_files(state)
schedule_poll(state.poll_interval)
{:noreply, state}
end
def handle_info({:DOWN, _ref, :process, pid, _reason}, state) do
{:noreply, %{state | subscribers: MapSet.delete(state.subscribers, pid)}}
end
# --- Private ---
defp schedule_poll(interval) do
Process.send_after(self(), :poll, interval)
end
defp workspace_root do
agent_dir =
try do
ExClaw.Config.get([:agent_dir])
rescue
_ -> nil
end
agent_dir || Path.join(System.get_env("HOME", "/tmp"), ".ex_claw/agent")
end
defp load_all_files(state) do
Enum.reduce(@file_atoms, state, fn {filename, name}, acc ->
path = Path.join(acc.root_dir, filename)
case read_file_meta(path) do
{:ok, meta, content} ->
acc
|> put_in([:file_meta, Access.key(name)], meta)
|> put_in([:cache, Access.key(name)], content)
:not_found ->
acc
end
end)
end
defp poll_files(state) do
# Check known files for changes
state = poll_known_files(state)
# Check for new files (e.g., files created since last poll)
state = check_new_files(state)
# Check daily memory directory for new files
poll_daily_memory(state)
end
defp poll_known_files(state) do
Enum.reduce(@file_atoms, state, fn {filename, name}, acc ->
path = Path.join(acc.root_dir, filename)
case read_file_meta(path) do
{:ok, meta, content} ->
old_meta = Map.get(acc.file_meta, name)
if old_meta && file_changed?(old_meta, meta) do
handle_file_change(acc, name, content, meta)
else
if old_meta == nil do
# File appeared — treat as new
handle_file_change(acc, name, content, meta)
else
acc
end
end
:not_found ->
if Map.has_key?(acc.file_meta, name) do
Logger.warning("Workspace file #{filename} was deleted, keeping cached version")
broadcast(acc.subscribers, {:file_deleted, name})
%{acc | file_meta: Map.delete(acc.file_meta, name)}
else
acc
end
end
end)
end
defp check_new_files(state) do
Enum.reduce(@file_atoms, state, fn {filename, name}, acc ->
if not Map.has_key?(acc.cache, name) do
path = Path.join(acc.root_dir, filename)
case read_file_meta(path) do
{:ok, meta, content} ->
Logger.info("Workspace file #{filename} detected")
acc
|> put_in([:file_meta, Access.key(name)], meta)
|> put_in([:cache, Access.key(name)], content)
|> tap(fn s -> broadcast(s.subscribers, {:file_changed, name, content}) end)
:not_found ->
acc
end
else
acc
end
end)
end
defp poll_daily_memory(state) do
ws_dir = Map.get(state, :context_dir) || ExClaw.Workspace.context_dir()
daily_dir = Path.join(ws_dir, "memory/daily")
if File.dir?(daily_dir) do
today = Date.to_iso8601(Date.utc_today())
today_file = Path.join(daily_dir, "#{today}.md")
case read_file_meta(today_file) do
{:ok, meta, content} ->
old_meta = Map.get(state.file_meta, :daily_memory)
if old_meta == nil || file_changed?(old_meta, meta) do
state
|> put_in([:file_meta, Access.key(:daily_memory)], meta)
|> put_in([:cache, Access.key(:daily_memory)], content)
|> tap(fn s -> broadcast(s.subscribers, {:file_changed, :daily_memory, content}) end)
else
state
end
:not_found ->
state
end
else
state
end
end
defp handle_file_change(state, name, content, meta) do
case Validator.validate(name, content) do
:ok ->
Logger.debug("Workspace file #{name} changed (valid)")
broadcast(state.subscribers, {:file_changed, name, content})
state
|> put_in([:file_meta, Access.key(name)], meta)
|> put_in([:cache, Access.key(name)], content)
{:warning, reason} ->
Logger.warning("Workspace file #{name} changed with warning: #{reason}")
broadcast(state.subscribers, {:file_changed, name, content})
state
|> put_in([:file_meta, Access.key(name)], meta)
|> put_in([:cache, Access.key(name)], content)
{:error, reason} ->
Logger.warning("Workspace file #{name} invalid: #{reason} — keeping cached version")
broadcast(state.subscribers, {:file_invalid, name, reason})
# Keep old cached content, but update meta to avoid re-checking
put_in(state, [:file_meta, Access.key(name)], meta)
end
end
defp read_file_meta(path) do
case File.stat(path) do
{:ok, %File.Stat{type: :regular, mtime: mtime, size: size}} ->
case File.read(path) do
{:ok, content} -> {:ok, %{mtime: mtime, size: size}, content}
{:error, _} -> :not_found
end
_ ->
:not_found
end
end
defp file_changed?(old_meta, new_meta) do
old_meta.mtime != new_meta.mtime || old_meta.size != new_meta.size
end
defp broadcast(subscribers, message) do
Enum.each(subscribers, fn pid ->
send(pid, message)
end)
end
end