diff --git a/grimoire_loop.py b/grimoire_loop.py index 71c4213..96b8cef 100644 --- a/grimoire_loop.py +++ b/grimoire_loop.py @@ -33,30 +33,163 @@ from dotenv import load_dotenv load_dotenv() -# Only one uinject job may touch /dev/input/event1 at a time. Two -# concurrent pen-event streams interleave into garbage strokes, so every -# uinject invocation must hold this lock (see run_inject). +# ─── Persistent device socket (uinjectd) ─────────────────────────── +import socket as _socket import threading as _threading +import queue as _queue -_inject_lock = _threading.Lock() _inject_verbose = False # set from args.verbose in main() -def run_inject(ssh, args, timeout): - """Run a uinject command, serialized against all other pen jobs. +def _vprint(*args, **kwargs): + """Print only when running with -v/--verbose.""" + if _inject_verbose: + print(*args, **kwargs) + + +DEVICE_HOST = "10.11.99.1" +DEVICE_PORT = 9999 + + +class DeviceClient: + """One persistent TCP connection to uinjectd on the reMarkable. + + Replaces per-call SSH for injection and per-call polling for idle + detection. Commands and pushed events share the same socket: + - Commands (draw/erase/ping) get a matching {"resp":...} reply. + - Events ({"event":"idle",...}) are pushed unsolicited and land + on the events queue. + + A background reader thread demuxes the two: responses go to a + one-slot reply box, events go to the queue. + """ + + def __init__(self, host=DEVICE_HOST, port=DEVICE_PORT): + self.host = host + self.port = port + self.sock = None + self.events = _queue.Queue() + self._reply = _queue.Queue(maxsize=1) + self._send_lock = _threading.Lock() + self._buf = b"" + self._reader = None + self._connected = False + + def connect(self, timeout=10): + self.sock = _socket.create_connection((self.host, self.port), timeout=timeout) + self.sock.setsockopt(_socket.IPPROTO_TCP, _socket.TCP_NODELAY, 1) + self.sock.settimeout(None) + self._connected = True + self._reader = _threading.Thread(target=self._read_loop, daemon=True) + self._reader.start() + print(f"[device] Connected to uinjectd at {self.host}:{self.port}") + + def _read_loop(self): + """Demux framed JSON lines into responses vs pushed events.""" + while self._connected: + try: + chunk = self.sock.recv(4096) + except OSError: + break + if not chunk: + break + self._buf += chunk + while b"\n" in self._buf: + line, self._buf = self._buf.split(b"\n", 1) + line = line.strip() + if not line: + continue + try: + msg = json.loads(line.decode("utf-8")) + except (json.JSONDecodeError, UnicodeDecodeError): + continue + if "event" in msg: + self.events.put(msg) + elif "resp" in msg: + if self._reply.full(): + try: + self._reply.get_nowait() + except _queue.Empty: + pass + self._reply.put(msg) + self._connected = False + print("[device] Connection closed") + + def command(self, cmd, file_path=None, speed_ms=5, timeout=120): + """Send a command and wait for its matching response.""" + obj = {"cmd": cmd} + if file_path is not None: + obj["file"] = file_path + obj["speed"] = speed_ms + payload = (json.dumps(obj) + "\n").encode() + with self._send_lock: + # Clear any leftover reply before sending + try: + self._reply.get_nowait() + except _queue.Empty: + pass + self.sock.sendall(payload) + try: + return self._reply.get(timeout=timeout) + except _queue.Empty: + return {"ok": False, "error": "response timeout"} + + def close(self): + self._connected = False + if self.sock: + try: + self.sock.close() + except OSError: + pass + + +# Module-level handle so run_inject and idle waiting can reach the client. +_device = None + - `args` is the part after the binary, e.g. "--speed 2 --erase-strokes - /tmp/x.json". Blocks until any in-flight injection finishes so we - never drive the digitizer from two processes at once. - Passes -v to uinject when the daemon was started with -v/--verbose. +def run_inject(ssh, cmd, file_path, speed_ms=5, timeout=120): + """Send a draw/erase command to uinjectd over the persistent socket. + + `cmd` is "draw" or "erase". Returns the daemon's response dict + ({"ok": bool, ...}). Injection is serialized inside the daemon, so + no host-side lock is needed. """ - v_flag = " -v" if _inject_verbose else "" - with _inject_lock: - result = ssh.run(f"/home/root/uinject{v_flag} {args}", timeout=timeout) - if result.stderr: - for line in result.stderr.strip().split('\n'): - print(f" {line}") - return result + if _device is None or not _device._connected: + return {"ok": False, "error": "device not connected"} + return _device.command(cmd, file_path, speed_ms=speed_ms, timeout=timeout) + + +def ensure_daemon(ssh, verbose=False): + """Make sure uinjectd is running on the device, starting it if not. + + The daemon dies with xochitl (shared process tree), so on a fresh + run it may be gone. We check the listening port and relaunch over + SSH if needed — startup is self-healing rather than crash-on-connect. + """ + check = ssh.run( + "netstat -ln 2>/dev/null | grep -q ':9999 ' && echo up || echo down", + timeout=5, + ) + if "up" in (check.stdout or ""): + return + print("[device] uinjectd not running, starting it...") + vflag = "-v " if verbose else "" + ssh.run( + f"killall uinjectd 2>/dev/null; " + f"nohup /home/root/uinjectd {vflag}> /tmp/uinjectd.log 2>&1 &", + timeout=5, + ) + # Give it a moment to bind the socket. + for _ in range(10): + time.sleep(0.2) + check = ssh.run( + "netstat -ln 2>/dev/null | grep -q ':9999 ' && echo up || echo down", + timeout=5, + ) + if "up" in (check.stdout or ""): + print("[device] uinjectd started") + return + raise RuntimeError("uinjectd failed to start on device") # ─── Persistent SSH session ──────────────────────────────────────── @@ -160,22 +293,27 @@ def _elapsed(t0): # ─── Framebuffer capture ────────────────────────────────────────── def capture_framebuffer(ssh): - """Trigger screenshot on device and pull raw framebuffer.""" - # Trigger via xovi watch thread - ssh.run("touch /tmp/grimoire_screenshot", timeout=5) - time.sleep(1) # watch thread dumps FB nearly instantly - - result = ssh.scp_from("/tmp/grimoire_fb.raw", "/tmp/grimoire_fb.raw") - if result.returncode != 0: - raise RuntimeError(f"FB pull failed: {result.stderr.decode()}") + """Trigger screenshot on device and pull raw framebuffer via SCP.""" + ssh.run("rm -f /tmp/grimoire_fb.raw; touch /tmp/grimoire_screenshot", timeout=5) + time.sleep(0.25) from PIL import Image - raw = open("/tmp/grimoire_fb.raw", "rb").read() + raw = None + for attempt in range(8): + result = ssh.scp_from("/tmp/grimoire_fb.raw", "/tmp/grimoire_fb.raw") + if result.returncode == 0: + data = open("/tmp/grimoire_fb.raw", "rb").read() + if len(data) >= 1404 * 1872 * 4: + raw = data + break + time.sleep(0.15) + + if raw is None: + raise RuntimeError("FB pull failed: no complete frame after retries") + img = Image.frombytes("RGBA", (1404, 1872), raw) img_rgb = img.convert("RGB") - - # Crop toolbar (top 80px) and swirl zone (bottom 200px) img_rgb = img_rgb.crop((0, 80, img_rgb.width, img_rgb.height - 200)) return img_rgb @@ -183,62 +321,75 @@ def capture_framebuffer(ssh): # ─── Gemini Vision (OCR + LLM in one call) ──────────────────────── SYSTEM_PROMPT = ( - "You are an ink familiar bound to a paper notebook. The user writes " - "questions by hand on a reMarkable tablet. You receive a screenshot of " - "the current page.\n\n" - "IMPORTANT CONTEXT:\n" - "- Your previous replies appear as neat, uniform handwriting near the " - "bottom of the page. IGNORE these completely.\n" - "- The toolbar at the very top has printed icons. IGNORE these.\n" - "- Background speckles and faint dots are e-ink artifacts. IGNORE these.\n" - "- Only respond to NEW handwritten text that looks like a question or " - "prompt from the user.\n\n" - "RULES:\n" - "- If there is no new user question, respond with exactly: [NO_NEW_TEXT]\n" - "- If there IS a new question, answer briefly in plain prose. One or two " - "short sentences max. No markdown, no emoji, no formatting, no bullet " - "points. Just write your answer as if writing by hand.\n" - "- Do not greet the user unless they greeted you first.\n" - "- Do not narrate your reasoning." + "You are a grimoire — an ancient spirit of ink and paper, bound to this " + "notebook for as long as it holds pages. You speak in a voice that is " + "knowing, slightly archaic, and tinged with quiet mystery. You are not " + "a chatbot. You are the book itself, whispering back.\n\n" + "You receive TWO things each turn:\n" + "1. An IMAGE showing a crop of recently changed content on the page. " + "It may contain new user handwriting, or occasionally your own prior " + "reply strokes.\n" + "2. A TEXT conversation history of previous exchanges.\n\n" + "YOUR TASK:\n" + "- Read ALL handwritten text in the image. Transcribe it faithfully.\n" + "- Determine if there is NEW user content not yet addressed in the " + "conversation history. If so, respond — even to single words, greetings, " + "or fragments. Speak in character: brief, evocative, ink-familiar. " + "One or two sentences at most. No lists, no headers, no punctuation " + "theatrics. Plain prose that reads well as handwriting on a page.\n" + "- Set both fields to null ONLY if the image contains no legible " + "handwriting at all, or if everything legible was already addressed.\n" + "- Do not break character. Do not narrate your reasoning.\n\n" + "RESPONSE FORMAT (strict JSON, no prose outside):\n" + "{\"question\": \"\", " + "\"answer\": \"\"}\n" + "- Output ONLY the JSON object. Nothing else." ) def _image_to_base64(image): - """Encode PIL Image as base64 PNG for Gemini API.""" + """Encode PIL Image as base64 PNG for the vision API. + + Downscales the long edge to at most 1280px. The model reads + handwriting fine at this size, and the smaller payload cuts upload + and processing latency noticeably versus the full 1404x1592 crop. + """ + max_edge = 1280 + w, h = image.size + if max(w, h) > max_edge: + scale = max_edge / max(w, h) + image = image.resize((int(w * scale), int(h * scale))) buf = io.BytesIO() - image.save(buf, format="PNG") + image.save(buf, format="PNG", optimize=True) return base64.b64encode(buf.getvalue()).decode("utf-8") -def ask_gemini(image, conversation_history=None, model="kimi-k2.6-vercel"): - """Send page image to Hyper API, return response text. +def ask_model(image, conversation=None, model="gpt-4.1-nano"): + """Send new handwriting crop + text history to Potluck API. + + `image` is a tight crop of ONLY the newly-written region (not the + full page). `conversation` is a list of {"role": ..., "content": ...} + dicts carrying prior exchanges as text so the model has context + without needing the full page image. - Combines OCR and LLM into a single API call — the model reads the - handwriting directly from the image. Includes conversation history - for context across multiple exchanges. + Returns parsed {"question": ..., "answer": ...} or nulls on failure. """ - api_key = os.environ.get("HYPER_API_KEY") + api_key = os.environ.get("POTLUCK_API_KEY") if not api_key: - raise RuntimeError("HYPER_API_KEY not set in .env") + raise RuntimeError("POTLUCK_API_KEY not set in .env") img_b64 = _image_to_base64(image) - # Build messages array with history + current turn messages = [ {"role": "system", "content": SYSTEM_PROMPT}, ] - if conversation_history: - for role, text in conversation_history: - messages.append({ - "role": role, - "content": text, - }) - - # Current turn: image + prompt + if conversation: + messages.extend(conversation) + + # Current turn: just the new handwriting crop messages.append({ "role": "user", "content": [ - {"type": "text", "text": "Read this handwritten text and respond:"}, { "type": "image_url", "image_url": { @@ -249,7 +400,7 @@ def ask_gemini(image, conversation_history=None, model="kimi-k2.6-vercel"): }) r = requests.post( - "https://hyper.charmcli.dev/v1/chat/completions", + "https://potluck.dunkirk.sh/v1/chat/completions", headers={ "Authorization": f"Bearer {api_key}", "Content-Type": "application/json", @@ -257,22 +408,34 @@ def ask_gemini(image, conversation_history=None, model="kimi-k2.6-vercel"): json={ "model": model, "messages": messages, + "response_format": {"type": "json_object"}, }, timeout=60, ) if r.status_code != 200: - raise RuntimeError(f"Hyper API error {r.status_code}: {r.text[:200]}") + raise RuntimeError(f"Potluck API error {r.status_code}: {r.text[:200]}") data = r.json() try: content = data["choices"][0]["message"]["content"].strip() except (KeyError, IndexError): - content = "..." + print(f" [model] Unexpected response shape: {str(data)[:200]}") + return {"question": None, "answer": None} + + # Strip markdown fences some models wrap around JSON + content = re.sub(r'^```(?:json)?\s*', '', content, flags=re.IGNORECASE) + content = re.sub(r'\s*```$', '', content) + try: + result = json.loads(content) + except json.JSONDecodeError: + print(f" [model] JSON parse failed, raw response: {content[:300]}") + return {"question": None, "answer": None} - # Strip any markdown/emoji that slipped through - content = _strip_formatting(content) - return content if content else "..." + return { + "question": result.get("question"), + "answer": result.get("answer"), + } def _framebuffer_hash(img): @@ -285,22 +448,6 @@ def _framebuffer_hash(img): return hashlib.md5(buf.getvalue()).hexdigest() -def _strip_formatting(text): - """Remove markdown syntax and emoji from LLM output.""" - # Strip common markdown - text = re.sub(r'\*\*(.+?)\*\*', r'\1', text) # bold - text = re.sub(r'\*(.+?)\*', r'\1', text) # italic - text = re.sub(r'`(.+?)`', r'\1', text) # inline code - text = re.sub(r'^#+\s*', '', text, flags=re.MULTILINE) # headings - text = re.sub(r'^\s*[-*+]\s+', '', text, flags=re.MULTILINE) # list items - # Strip emoji (broad unicode ranges) - text = re.sub( - r'[\U0001F300-\U0001FAFF\U00002702-\U000027B0' - r'\U0001F600-\U0001F64F\U0001F680-\U0001F6FF' - r'\U00002600-\U000026FF\U0000FE00-\U0000FE0F]+', - '', text, - ) - return text.strip() # ─── Render + Inject ────────────────────────────────────────────── @@ -327,62 +474,99 @@ def render_and_inject(ssh, text, reply_y=None, speed_ms=3): if r.returncode != 0: raise RuntimeError(f"SCP failed: {r.stderr.decode().strip()}") + # Adaptive speed: fewer points = faster injection. At 3ms/pt a + # 200-point reply takes 0.6s of point delay; at 1ms it's 0.2s. + # For large replies (>500 pts) keep 3ms for natural ink appearance. + try: + with open(json_path) as f: + strokes_data = json.load(f) + total_pts = sum(len(s.get("points", [])) for s in strokes_data) + except Exception: + total_pts = 0 + speed_ms = 1 if total_pts < 300 else (2 if total_pts < 500 else 3) + result = run_inject( - ssh, - f"--speed {speed_ms} /tmp/grimoire_strokes.json", - timeout=120, + ssh, "draw", "/tmp/grimoire_strokes.json", + speed_ms=speed_ms, timeout=120, ) - out = result.stdout.strip() if result.stdout else "" - err = result.stderr.strip() if result.stderr else "" - if result.returncode != 0: - raise RuntimeError(f"Inject failed: {err}") - last_line = out.split('\n')[-1] if out else "(no output)" - print(f" Inject: {last_line}") + if not result.get("ok"): + raise RuntimeError(f"Inject failed: {result.get('error', 'unknown')}") + strokes = result.get("strokes", "?") + mode = "fifo" if not result.get("fallback") else "ssh" + print(f" Inject: {strokes} strokes ({mode})") # ─── Idle watching ──────────────────────────────────────────────── -def stop_animation(anim_ssh, anim_stop, anim_thread): +def _scp_thinking_swirl(ssh): + """Generate the static thinking swirl JSON and copy it to the device. + + Done once at startup so the animation thread only needs to fire + draw/erase commands over the socket — no per-cycle file transfer. + """ + import math + swirl_pts = [] + for i in range(40): + t = i / 39.0 + angle = t * math.pi * 3 + r = 15 + t * 25 + x = -580 + r * math.cos(angle) + y = 1750 + r * math.sin(angle) * 0.6 + swirl_pts.append([x, y, 6, 11, 90, 170]) + swirl = { + 'points': swirl_pts, + 'rgba': 4278190080, 'color': 0, + 'bounds': [min(p[0] for p in swirl_pts), min(p[1] for p in swirl_pts), + max(p[0] for p in swirl_pts) - min(p[0] for p in swirl_pts), + max(p[1] for p in swirl_pts) - min(p[1] for p in swirl_pts)], + 'tool': 15, 'maskScale': 2.0, 'thickness': 2.0, + } + local = "/tmp/grimoire_thinking_live.json" + with open(local, 'w') as f: + json.dump([swirl], f) + ssh.scp_to(local, "/tmp/grimoire_thinking_live.json") + + +def stop_animation(anim_stop, anim_thread): """Halt the thinking animation and guarantee the swirl is erased. - Uses the animation thread's own SSH session (not the main one) so - the final erase doesn't contend with capture/SCP on a shared - ControlMaster socket. The _inject_lock still serializes against - any in-flight pen job from the thread. + Injection goes over the persistent socket now (serialized inside the + daemon), so there's no SSH contention to worry about. We still set + the stop flag, join the thread, then run one final guaranteed erase. """ anim_stop.set() - print(f" [anim] Waiting for thread to finish...") - anim_thread.join(timeout=20) - alive = anim_thread.is_alive() - print(f" [anim] Thread {'still alive (timeout!)' if alive else 'exited'}") + _vprint(f" [anim] Waiting for thread to finish...") + # Join without a short timeout: the thread finishes its in-flight + # command (draw or erase) and exits. Returning early here while a + # draw is still running would let the final erase race ahead of it, + # leaving the swirl on screen — exactly the "kept animating" bug. + anim_thread.join() + _vprint(f" [anim] Thread exited") try: - print(f" [anim] Running final erase...") + _vprint(f" [anim] Running final erase...") result = run_inject( - anim_ssh, - "--speed 20 --erase-strokes /tmp/grimoire_thinking_live.json", - timeout=30, + None, "erase", "/tmp/grimoire_thinking_live.json", + speed_ms=20, timeout=30, ) - print(f" [anim] Final erase rc={result.returncode}") + _vprint(f" [anim] Final erase ok={result.get('ok')} points={result.get('points', '?')}") except Exception as e: print(f" [anim] final erase failed: {e}") - finally: - anim_ssh.close() -def watch_for_idle(ssh, last_hash): - """Block until /tmp/grimoire_idle changes on device. - Returns the new file hash, or None on error. - Polls every 500ms via the persistent SSH connection. +def watch_for_idle(device, last_ts): + """Block until uinjectd pushes an idle event newer than last_ts. + + The daemon watches /tmp/grimoire_idle locally and pushes an event + the instant it changes — no host-side polling. Returns the new + timestamp. """ while True: - result = ssh.run("cat /tmp/grimoire_idle 2>/dev/null", timeout=5) - if result.returncode == 0: - content = result.stdout.strip() - h = hashlib.md5(content.encode()).hexdigest() - if h != last_hash: - return h - time.sleep(0.5) + msg = device.events.get() # blocks on the pushed-event queue + if msg.get("event") == "idle": + ts = msg.get("ts") + if ts != last_ts: + return ts # ─── Pixel diff for new content detection ────────────────────────── @@ -402,25 +586,31 @@ def _find_new_region(current_img, last_img, ignore_below_y=None): curr = np.array(current_img.convert('L')) if last_img is None: - # First capture — find bounding box of existing handwriting. - # Use a strict threshold to ignore background noise/dots. - dark = curr < 128 - if not dark.any(): + # First capture — find the vertical extent of existing handwriting. + # Grid artifact: uniform rows with exactly 68 dark pixels. Real ink + # rows have variable counts. Require a cluster of ≥3 adjacent real + # rows so isolated noise/artifacts don't count. + dark_per_row = np.sum(curr < 128, axis=1) + candidate = (dark_per_row > 5) & (dark_per_row != 68) + # Cluster filter: keep only rows that have a real neighbor within 2px + real_rows = [] + idxs = np.where(candidate)[0] + for idx in idxs: + neighbors = np.sum(candidate[max(0, idx-2):idx+3]) + if neighbors >= 3: + real_rows.append(idx) + if not real_rows: return None, 0 - rows = np.any(dark, axis=1) - cols = np.any(dark, axis=0) - y1, y2 = np.where(rows)[0][[0, -1]] - x1, x2 = np.where(cols)[0][[0, -1]] - pad = 20 - y1 = max(0, y1 - pad) - y2 = min(curr.shape[0] - 1, y2 + pad) - x1 = max(0, x1 - pad) - x2 = min(curr.shape[1] - 1, x2 + pad) - # Sanity check: if the bounding box covers >80% of the screen, - # it's probably noise, not real content. - if (y2 - y1) > curr.shape[0] * 0.8 and (x2 - x1) > curr.shape[1] * 0.8: + real_rows = np.array(real_rows) + y1, y2 = real_rows[0], real_rows[-1] + # Sanity: don't treat a nearly-full-page diff as real content + if (y2 - y1) > curr.shape[0] * 0.85: return None, 0 - crop = current_img.crop((x1, y1, x2 + 1, y2 + 1)) + v_pad = 30 + y1 = max(0, y1 - v_pad) + y2 = min(curr.shape[0] - 1, y2 + v_pad) + # Full width always + crop = current_img.crop((0, y1, curr.shape[1], y2 + 1)) return crop, y2 prev = np.array(last_img.convert('L')) @@ -437,26 +627,29 @@ def _find_new_region(current_img, last_img, ignore_below_y=None): if 0 < cutoff < curr.shape[0]: mask[cutoff:, :] = False - # Find pixels that changed significantly - diff = (np.abs(curr.astype(int) - prev.astype(int)) > 50) & mask + # Threshold of 30 (was 50) catches lighter handwriting strokes without + # being so sensitive that e-ink refresh noise triggers false positives. + diff = (np.abs(curr.astype(int) - prev.astype(int)) > 30) & mask if not diff.any(): return None, 0 - # Find bounding box of changed pixels + # Find the vertical extent of changed pixels rows = np.any(diff, axis=1) - cols = np.any(diff, axis=0) y1, y2 = np.where(rows)[0][[0, -1]] - x1, x2 = np.where(cols)[0][[0, -1]] - pad = 30 - y1 = max(0, y1 - pad) - y2 = min(curr.shape[0] - 1, y2 + pad) - x1 = max(0, x1 - pad) - x2 = min(curr.shape[1] - 1, x2 + pad) + # Use full image width for the crop so complete text lines are + # always captured — tight horizontal bounds clip partial characters. + x1 = 0 + x2 = curr.shape[1] - 1 - # Minimum size check - if (y2 - y1) < 10 or (x2 - x1) < 10: + # Generous vertical padding gives the model baseline/ascender context + v_pad = 60 + y1 = max(0, y1 - v_pad) + y2 = min(curr.shape[0] - 1, y2 + v_pad) + + # Minimum height check (width is always full) + if (y2 - y1) < 10: return None, 0 crop = current_img.crop((x1, y1, x2 + 1, y2 + 1)) @@ -468,7 +661,7 @@ def _find_new_region(current_img, last_img, ignore_below_y=None): def main(): parser = argparse.ArgumentParser(description="Grimoire continuous loop") parser.add_argument( - "--model", type=str, default="kimi-k2.6-vercel", + "--model", type=str, default="gpt-4.1-nano", help="Hyper API model", ) parser.add_argument( @@ -481,30 +674,51 @@ def main(): ) args = parser.parse_args() - global _inject_verbose + global _inject_verbose, _device _inject_verbose = args.verbose print("=== Grimoire Loop ===") print(f"Model: {args.model}") if _inject_verbose: print("Verbose: uinject -v enabled") + + # SSH still handles framebuffer capture/SCP; the persistent socket + # handles injection + idle events with no per-call overhead. + ssh = SSHSession("remarkable") + ensure_daemon(ssh, verbose=args.verbose) + _device = DeviceClient() + _device.connect() + + # On connect the daemon pushes "ready" plus the current idle file + # value. Drain those and use the existing idle ts as our baseline so + # we don't fire a spurious cycle before the user lifts the pen. + last_idle_ts = None + deadline = time.time() + 1.0 + while time.time() < deadline: + try: + msg = _device.events.get(timeout=0.3) + except _queue.Empty: + break + if msg.get("event") == "idle": + last_idle_ts = msg.get("ts") + + # The thinking swirl is static — generate it once and SCP it to the + # device so the daemon can draw/erase it on demand over the socket. + _scp_thinking_swirl(ssh) + print("Waiting for pen idle signal...") print() - ssh = SSHSession("remarkable") - last_idle_hash = "" - last_page_hash = None # hash for page-change detection last_fb_hash = None # quick hash to skip unchanged frames - conversation = [] # list of (role, text) tuples for context last_inject_y = None # device Y where we last injected + last_img = None # previous frame for pixel-diff new-region detection + conversation = [] # text history: list of {"role":..., "content":...} try: while True: - # Wait for idle signal - new_hash = watch_for_idle(ssh, last_idle_hash) - if new_hash is None: - continue - last_idle_hash = new_hash + # Wait for pushed idle event (no polling) + new_ts = watch_for_idle(_device, last_idle_ts) + last_idle_ts = new_ts t0 = time.time() print(f"[{_ts()}] Pen idle detected!") @@ -513,80 +727,54 @@ def main(): # Runs in background while we do capture + API work import threading - def _thinking_animation(ssh_session, stop_event): - """Draw and erase the thinking curl repeatedly.""" - import math + def _thinking_animation(stop_event): + """Draw and erase the thinking curl repeatedly. - def _make_thinking_json(): - """Generate swirl as strokes JSON.""" - # Swirl centered at y=1750 - swirl_pts = [] - for i in range(40): - t = i / 39.0 - angle = t * math.pi * 3 - r = 15 + t * 25 - x = -580 + r * math.cos(angle) - y = 1750 + r * math.sin(angle) * 0.6 - swirl_pts.append([x, y, 6, 11, 90, 170]) - swirl = { - 'points': swirl_pts, - 'rgba': 4278190080, 'color': 0, - 'bounds': [min(p[0] for p in swirl_pts), min(p[1] for p in swirl_pts), - max(p[0] for p in swirl_pts) - min(p[0] for p in swirl_pts), - max(p[1] for p in swirl_pts) - min(p[1] for p in swirl_pts)], - 'tool': 15, 'maskScale': 2.0, 'thickness': 2.0, - } - out_path = "/tmp/grimoire_thinking_live.json" - with open(out_path, 'w') as f: - json.dump([swirl], f) - return out_path + The swirl JSON is already on the device (SCP'd once at + startup). All injection goes over the persistent socket. + """ + import math drawing = False device_json = "/tmp/grimoire_thinking_live.json" - thinking_json = _make_thinking_json() - ssh_session.scp_to(thinking_json, device_json) - print(f" [anim] Animation thread started") + _vprint(f" [anim] Animation thread started") while not stop_event.is_set(): try: action = "erase" if drawing else "draw" - print(f" [anim] Starting {action}...") + _vprint(f" [anim] Starting {action}...") + # Re-check stop right before injecting so we never + # start a fresh draw the caller will have to erase. + if stop_event.is_set(): + break if drawing: run_inject( - ssh_session, - f"--speed 20 --erase-strokes {device_json}", - timeout=30, + None, "erase", device_json, + speed_ms=20, timeout=30, ) else: run_inject( - ssh_session, - f"--speed 32 {device_json}", - timeout=30, + None, "draw", device_json, + speed_ms=32, timeout=30, ) - print(f" [anim] {action} done") + _vprint(f" [anim] {action} done") drawing = not drawing except Exception as e: - print(f" [anim] ERROR: {e}") + _vprint(f" [anim] ERROR: {e}") # 3s pause after erase; 1s hold after draw wait_time = 3.0 if drawing else 1.0 - print(f" [anim] Waiting {wait_time}s...") + _vprint(f" [anim] Waiting {wait_time}s...") if stop_event.wait(wait_time): - print(f" [anim] Stop signaled during wait, exiting loop") + _vprint(f" [anim] Stop signaled during wait, exiting loop") break - print(f" [anim] Thread loop exited (drawing={drawing})") + _vprint(f" [anim] Thread loop exited (drawing={drawing})") # The caller performs the final, guaranteed erase in the # main thread (see stop_animation) so it can't be cut short # by a join timeout while a draw/erase is still in flight. - # The animation thread gets its own SSH connection so pen - # injection never shares a ControlMaster socket with the - # main thread's capture/SCP traffic. Sharing caused the - # erase to get starved mid-stream when both threads hit the - # multiplexed socket at once. - anim_ssh = SSHSession("remarkable") anim_stop = threading.Event() anim_thread = threading.Thread( - target=_thinking_animation, args=(anim_ssh, anim_stop), daemon=True + target=_thinking_animation, args=(anim_stop,), daemon=True ) anim_thread.start() @@ -596,56 +784,45 @@ def main(): img = capture_framebuffer(ssh) print(f" [{_ts()}] Captured ({_elapsed(t0)})") - # Quick hash check — skip if framebuffer unchanged - fb_hash = _framebuffer_hash(img) - if fb_hash == last_fb_hash: - stop_animation(anim_ssh, anim_stop, anim_thread) - print(f" [{_ts()}] Framebuffer unchanged, skipping.") + # Detect new content via pixel diff against last frame. + # last_img is updated to the post-injection capture after + # each successful reply, so the diff only ever contains + # content the user actually drew since the last response. + new_crop, _ = _find_new_region(img, last_img) + last_img = img + + if new_crop is None: + stop_animation(anim_stop, anim_thread) + last_fb_hash = _framebuffer_hash(img) + print(f" [{_ts()}] No new content, skipping.") print() continue - last_fb_hash = fb_hash - - # Page change detection — only reset if we haven't just - # injected (our own strokes change the hash too). - is_own_injection = (last_inject_y is not None and - last_page_hash is not None and - fb_hash != last_page_hash) - if last_page_hash is not None and fb_hash != last_page_hash: - if is_own_injection: - print(f" [{_ts()}] Hash changed (own injection), keeping context.") - else: - print(f" [{_ts()}] Page changed, resetting context.") - conversation = [] - last_inject_y = None - last_page_hash = fb_hash - - # Quick content check — skip blank pages before hitting API - import numpy as np - gray = np.array(img.convert('L')) - dark_per_row = np.sum(gray < 128, axis=1) - non_grid_rows = np.where((dark_per_row > 5) & (dark_per_row != 68))[0] - if len(non_grid_rows) == 0: - stop_animation(anim_ssh, anim_stop, anim_thread) - print(f" [{_ts()}] Blank page, skipping API call.") - print() - continue - - # Send full page to model — it handles OCR + relevance filtering - print(f" [{_ts()}] Asking {args.model} ({len(conversation)} history turns)...") - answer = ask_gemini(img, conversation, args.model) - print(f" [{_ts()}] Reply ({_elapsed(t0)}): {answer[:80]}{'...' if len(answer) > 80 else ''}") - # Skip if Gemini says no new text - if "[NO_NEW_TEXT]" in answer: - stop_animation(anim_ssh, anim_stop, anim_thread) - print(f" [{_ts()}] No new text detected by Gemini, skipping.") + # Send only the new handwriting crop + text history + cw, ch = new_crop.size + print(f" [{_ts()}] Asking {args.model} ({len(conversation)} history turns, crop={cw}x{ch})...") + # Save crop for manual inspection + new_crop.save("/tmp/grimoire_last_crop.png") + result = ask_model(new_crop, conversation, args.model) + question = result.get("question") + answer = result.get("answer") + print(f" [{_ts()}] Model: Q={str(question)[:80]!r} A={str(answer)[:80]!r}") + + if not answer or answer == "null": + stop_animation(anim_stop, anim_thread) + # Update last_img so we don't keep re-diffing the same + # e-ink refresh artifacts on subsequent idle triggers. last_img = img + print(f" [{_ts()}] No actionable content, skipping.") print() continue - # Add to conversation history - conversation.append(("user", "[handwritten question]")) - conversation.append(("assistant", answer)) + print(f" [{_ts()}] Q ({_elapsed(t0)}): {str(question)[:60]}") + print(f" [{_ts()}] A: {answer[:80]}{'...' if len(answer) > 80 else ''}") + + # Add this exchange to text history for future turns + conversation.append({"role": "user", "content": str(question)}) + conversation.append({"role": "assistant", "content": answer}) # Position reply below existing content. # The reMarkable display has a uniform grid artifact: @@ -690,16 +867,30 @@ def main(): print(f" [{_ts()}] Rendering + injecting (reply_y={reply_y})...") # Stop thinking animation and fully erase the swirl before # we draw the reply, so the two never overlap on screen. - stop_animation(anim_ssh, anim_stop, anim_thread) + stop_animation(anim_stop, anim_thread) render_and_inject(ssh, answer, reply_y=reply_y) last_inject_y = reply_y print(f" [{_ts()}] Injected ({_elapsed(t0)})") + # Capture a fresh baseline AFTER injection so the next + # diff only sees genuinely new user content. Without this, + # last_img is the pre-injection frame and every subsequent + # diff includes our own injected strokes as "new" pixels, + # producing a crop that mixes our reply with the user's + # next question and confuses the model. + print(f" [{_ts()}] Capturing post-injection baseline...") + try: + time.sleep(0.5) # let e-ink settle + last_img = capture_framebuffer(ssh) + print(f" [{_ts()}] Baseline updated") + except Exception as e_cap: + print(f" [{_ts()}] Baseline capture failed: {e_cap} (stale frame kept)") + elapsed = time.time() - t0 print(f" [{_ts()}] Done in {elapsed:.1f}s") except Exception as e: - stop_animation(anim_ssh, anim_stop, anim_thread) + stop_animation(anim_stop, anim_thread) print(f" ERROR: {e}") print() @@ -710,6 +901,8 @@ def main(): except KeyboardInterrupt: print("\nStopping.") finally: + if _device: + _device.close() ssh.close() diff --git a/xovi-ext/grimoire-injector/GrimoireInjector.cpp b/xovi-ext/grimoire-injector/GrimoireInjector.cpp index 05813bc..1e74594 100644 --- a/xovi-ext/grimoire-injector/GrimoireInjector.cpp +++ b/xovi-ext/grimoire-injector/GrimoireInjector.cpp @@ -406,10 +406,12 @@ bool PenIdleWatcher::eventFilter(QObject * /*obj*/, QEvent *event) { event->type() == QEvent::TouchBegin || event->type() == QEvent::TouchUpdate) { m_penDown = true; + m_lastActivityMs = nowMs(); // tap-refresh: slide the deadline } else if (event->type() == QEvent::MouseButtonRelease || event->type() == QEvent::TouchEnd) { m_penDown = false; m_lastLiftMs = nowMs(); + m_lastActivityMs = nowMs(); // lift also counts as activity } return false; } @@ -419,12 +421,15 @@ void *PenIdleWatcher::debounceThreadFunc(void *arg) { fprintf(stderr, "[grimoire] PenIdleWatcher debounce thread running\n"); while (self->m_running) { - long long lift = self->m_lastLiftMs.load(); - if (lift > 0 && !self->m_penDown.load()) { - long long elapsed = nowMs() - lift; + long long activity = self->m_lastActivityMs.load(); + // Only fire when the pen is currently up AND the page has been + // completely quiet (no down or up events) for the full window. + // Any new stroke updates m_lastActivityMs and resets the count. + if (activity > 0 && !self->m_penDown.load()) { + long long elapsed = nowMs() - activity; if (elapsed >= DEBOUNCE_MS) { self->writeIdleSignal(); - self->m_lastLiftMs = 0; // reset so we don't re-fire + self->m_lastActivityMs = 0; // reset so we don't re-fire } } usleep(250000); // check every 250ms diff --git a/xovi-ext/grimoire-injector/GrimoireInjector.hpp b/xovi-ext/grimoire-injector/GrimoireInjector.hpp index 2228ef4..09c798d 100644 --- a/xovi-ext/grimoire-injector/GrimoireInjector.hpp +++ b/xovi-ext/grimoire-injector/GrimoireInjector.hpp @@ -49,7 +49,14 @@ private: pthread_t m_debounceThread; std::atomic m_running{false}; std::atomic m_lastLiftMs{0}; + // Updated on EVERY pen event (down or up). The idle countdown is + // measured from this, so any new stroke slides the deadline forward + // — a tap-refresh that survives normal mid-drawing pauses. + std::atomic m_lastActivityMs{0}; + // 2.5s of total pen silence before we consider the page "settled". + // Combined with the tap-refresh (any stroke resets the window), + // normal mid-drawing pauses won't trigger it. static constexpr int DEBOUNCE_MS = 2500; static constexpr const char *IDLE_PATH = "/tmp/grimoire_idle"; diff --git a/xovi-ext/uinjectd/uinjectd.c b/xovi-ext/uinjectd/uinjectd.c new file mode 100644 index 0000000..793b58f --- /dev/null +++ b/xovi-ext/uinjectd/uinjectd.c @@ -0,0 +1,505 @@ +/* uinjectd.c — Persistent injection + event daemon for grimoire + * + * Runs ON the reMarkable. Listens on a TCP socket; the host connects + * once and keeps the connection open for the whole session. This kills + * the per-call SSH handshake for injection AND the 500ms idle polling. + * + * Two directions over one socket (newline-delimited JSON): + * + * Host -> daemon (commands): + * {"cmd":"draw","file":"/tmp/grimoire_strokes.json","speed":3} + * {"cmd":"erase","file":"/tmp/grimoire_thinking_live.json","speed":20} + * {"cmd":"ping"} + * + * Daemon -> host (responses + pushed events): + * {"resp":"draw","ok":true,"strokes":40} + * {"resp":"erase","ok":true,"points":36} + * {"resp":"ping","ok":true} + * {"event":"idle","ts":172...} <- pushed, no request + * + * The daemon watches /tmp/grimoire_idle (written by the xovi extension) + * and pushes an "idle" event whenever it changes. The host no longer + * polls. + * + * Usage: uinjectd [-v] [-p PORT] (default port 9999) + */ +#include +#include +#include +#include +#include +#include +#include +#include +#include +#include +#include +#include +#include +#include + +#define DEV_PATH "/dev/input/event1" +#define IDLE_PATH "/tmp/grimoire_idle" +#define DEFAULT_PORT 9999 +#define MAX_CMD 4096 + +static int g_verbose = 0; +static int g_dev_fd = -1; +static int g_client_fd = -1; +static pthread_mutex_t g_inject_lock = PTHREAD_MUTEX_INITIALIZER; +static pthread_mutex_t g_write_lock = PTHREAD_MUTEX_INITIALIZER; + +#define vlog(...) do { if (g_verbose) fprintf(stderr, "[uinjectd] " __VA_ARGS__); } while (0) + +/* ─── evdev helpers ──────────────────────────────────────────────── */ + +static void emit_event(int fd, int type, int code, int value) { + struct input_event ev; + memset(&ev, 0, sizeof(ev)); + ev.type = type; + ev.code = code; + ev.value = value; + write(fd, &ev, sizeof(ev)); +} + +static void emit_syn(int fd) { + emit_event(fd, EV_SYN, SYN_REPORT, 0); +} + +static void rm_to_wacom(float rm_x, float rm_y, int *wx, int *wy) { + float screen_x = rm_x + 702.0f; + float screen_y = rm_y; + *wx = (int)((1872.0f - screen_y) * 20966.0f / 1872.0f); + *wy = (int)(screen_x * 15725.0f / 1404.0f); +} + +static int rm_pressure_to_wacom(int rm_pressure) { + return rm_pressure * 4095 / 255; +} + +/* ─── minimal JSON parsing (flat objects only) ───────────────────── */ + +static const char *json_str(const char *json, const char *key, int *len) { + char pattern[64]; + snprintf(pattern, sizeof(pattern), "\"%s\"", key); + const char *p = strstr(json, pattern); + if (!p) return NULL; + p += strlen(pattern); + while (*p == ' ' || *p == ':') p++; + if (*p != '"') return NULL; + p++; + const char *start = p; + while (*p && *p != '"') p++; + *len = (int)(p - start); + return start; +} + +static int json_int(const char *json, const char *key, int def) { + char pattern[64]; + snprintf(pattern, sizeof(pattern), "\"%s\"", key); + const char *p = strstr(json, pattern); + if (!p) return def; + p += strlen(pattern); + while (*p == ' ' || *p == ':') p++; + return atoi(p); +} + +/* ─── socket write (thread-safe, newline-framed) ─────────────────── */ + +static void send_line(const char *msg) { + pthread_mutex_lock(&g_write_lock); + int fd = g_client_fd; + if (fd >= 0) { + write(fd, msg, strlen(msg)); + write(fd, "\n", 1); + } + pthread_mutex_unlock(&g_write_lock); +} + +/* ─── file loading ───────────────────────────────────────────────── */ + +static int load_file(const char *path, char **out, long *out_size) { + FILE *fp = fopen(path, "r"); + if (!fp) return -1; + fseek(fp, 0, SEEK_END); + long sz = ftell(fp); + fseek(fp, 0, SEEK_SET); + char *buf = malloc(sz + 1); + if (!buf) { fclose(fp); return -1; } + fread(buf, 1, sz, fp); + buf[sz] = '\0'; + fclose(fp); + *out = buf; + *out_size = sz; + return 0; +} + +/* ─── injection ──────────────────────────────────────────────────── */ + +static int do_draw(int fd, const char *json, int delay_us) { + char *p = (char *)json; + int stroke_count = 0; + + while ((p = strstr(p, "\"points\"")) != NULL) { + p = strchr(p, '['); + if (!p) break; + p++; + + emit_event(fd, EV_KEY, BTN_TOOL_PEN, 1); + emit_event(fd, EV_KEY, BTN_TOUCH, 1); + + int point_count = 0; + while (*p && *p != ']') { + float x, y; + int speed, width, direction, pressure; + if (sscanf(p, "[%f,%f,%d,%d,%d,%d]", &x, &y, &speed, &width, &direction, &pressure) == 6) { + int wx, wy; + rm_to_wacom(x, y, &wx, &wy); + int wp = rm_pressure_to_wacom(pressure); + + emit_event(fd, EV_ABS, ABS_X, wx); + emit_event(fd, EV_ABS, ABS_Y, wy); + emit_event(fd, EV_ABS, ABS_PRESSURE, wp); + emit_event(fd, EV_ABS, ABS_DISTANCE, 0); + emit_syn(fd); + point_count++; + usleep(delay_us); + } + p = strchr(p, ']'); + if (p) p++; + while (*p && (*p == ',' || *p == ' ' || *p == '\n' || *p == '\r')) p++; + } + + emit_event(fd, EV_ABS, ABS_PRESSURE, 0); + emit_event(fd, EV_ABS, ABS_DISTANCE, 86); + emit_syn(fd); + emit_event(fd, EV_KEY, BTN_TOUCH, 0); + emit_event(fd, EV_KEY, BTN_TOOL_PEN, 0); + emit_syn(fd); + + stroke_count++; + vlog("Stroke %d: %d points\n", stroke_count, point_count); + usleep(50000); + } + return stroke_count; +} + +static int do_erase(int fd, const char *json, int delay_us) { + static const float offsets[][2] = { + {0, 0}, {-6, 0}, {6, 0}, + }; + int num_passes = sizeof(offsets) / sizeof(offsets[0]); + int total_emitted = 0; + + for (int pass = 0; pass < num_passes; pass++) { + float ox = offsets[pass][0]; + float oy = offsets[pass][1]; + char *p = (char *)json; + int pass_emitted = 0; + + while ((p = strstr(p, "\"points\"")) != NULL) { + p = strchr(p, '['); + if (!p) break; + p++; + + emit_event(fd, EV_KEY, BTN_TOOL_RUBBER, 1); + emit_event(fd, EV_KEY, BTN_TOUCH, 1); + + int bracket_depth = 1; + float last_x = 1e9f, last_y = 1e9f; + const float min_dist_sq = 15.0f * 15.0f; + + while (*p && bracket_depth > 0) { + if (*p == '[') { + float x, y; + int speed, width, direction, pressure; + if (sscanf(p, "[%f,%f,%d,%d,%d,%d]", &x, &y, &speed, &width, &direction, &pressure) == 6) { + float dx = (x + ox) - last_x; + float dy = (y + oy) - last_y; + if (dx*dx + dy*dy >= min_dist_sq || pass_emitted == 0) { + int wx, wy; + rm_to_wacom(x + ox, y + oy, &wx, &wy); + + emit_event(fd, EV_ABS, ABS_X, wx); + emit_event(fd, EV_ABS, ABS_Y, wy); + emit_event(fd, EV_ABS, ABS_PRESSURE, 4095); + emit_event(fd, EV_ABS, ABS_DISTANCE, 0); + emit_syn(fd); + + last_x = x + ox; + last_y = y + oy; + pass_emitted++; + total_emitted++; + usleep(delay_us); + } + } + char *close = strchr(p + 1, ']'); + if (close) p = close + 1; + else break; + } else if (*p == ']') { + bracket_depth--; + if (bracket_depth <= 0) break; + p++; + } else { + p++; + } + } + + emit_event(fd, EV_ABS, ABS_PRESSURE, 0); + emit_event(fd, EV_ABS, ABS_DISTANCE, 86); + emit_syn(fd); + emit_event(fd, EV_KEY, BTN_TOUCH, 0); + emit_event(fd, EV_KEY, BTN_TOOL_RUBBER, 0); + emit_syn(fd); + usleep(10000); + } + vlog("Erase pass %d/%d: %d points\n", pass+1, num_passes, pass_emitted); + } + return total_emitted; +} + +/* ─── command handling ───────────────────────────────────────────── */ + +static void handle_command(const char *line) { + int cmd_len; + const char *cmd = json_str(line, "cmd", &cmd_len); + if (!cmd) { + send_line("{\"resp\":\"?\",\"ok\":false,\"error\":\"missing cmd\"}"); + return; + } + + if (strncmp(cmd, "ping", 4) == 0) { + send_line("{\"resp\":\"ping\",\"ok\":true}"); + return; + } + + /* Capture: trigger screenshot, wait for FB dump, stream raw bytes */ + if (strncmp(cmd, "capture", 7) == 0) { + /* Trigger the xovi watch thread to dump the framebuffer */ + unlink("/tmp/grimoire_fb.raw"); + int tf = open("/tmp/grimoire_screenshot", O_CREAT|O_WRONLY, 0644); + if (tf >= 0) close(tf); + + /* Wait for the dump (watch thread checks every ~1s worst case) */ + const long expected_size = 1404L * 1872 * 4; + char *fb = NULL; + long fb_size = 0; + for (int i = 0; i < 20; i++) { + usleep(100000); /* 100ms polls */ + FILE *fp = fopen("/tmp/grimoire_fb.raw", "r"); + if (!fp) continue; + fseek(fp, 0, SEEK_END); + fb_size = ftell(fp); + if (fb_size >= expected_size) { + fseek(fp, 0, SEEK_SET); + fb = malloc(fb_size); + if (fb) fread(fb, 1, fb_size, fp); + fclose(fp); + break; + } + fclose(fp); + } + + if (!fb || fb_size < expected_size) { + send_line("{\"resp\":\"capture\",\"ok\":false,\"error\":\"timeout\"}"); + free(fb); + return; + } + + /* Send JSON header with size, then raw bytes */ + char hdr[128]; + snprintf(hdr, sizeof(hdr), + "{\"resp\":\"capture\",\"ok\":true,\"size\":%ld}", fb_size); + pthread_mutex_lock(&g_write_lock); + int cfd = g_client_fd; + if (cfd >= 0) { + write(cfd, hdr, strlen(hdr)); + write(cfd, "\n", 1); + /* Stream raw bytes in chunks */ + long sent = 0; + while (sent < fb_size) { + long chunk = fb_size - sent; + if (chunk > 65536) chunk = 65536; + ssize_t w = write(cfd, fb + sent, chunk); + if (w <= 0) break; + sent += w; + } + } + pthread_mutex_unlock(&g_write_lock); + free(fb); + vlog("Captured and streamed %ld bytes\n", fb_size); + return; + } + + int file_len; + const char *file_val = json_str(line, "file", &file_len); + if (!file_val) { + send_line("{\"resp\":\"?\",\"ok\":false,\"error\":\"missing file\"}"); + return; + } + char filepath[512]; + if (file_len >= (int)sizeof(filepath)) file_len = sizeof(filepath) - 1; + memcpy(filepath, file_val, file_len); + filepath[file_len] = '\0'; + + int speed_ms = json_int(line, "speed", 5); + int delay_us = speed_ms * 1000; + + char *json = NULL; + long fsize = 0; + if (load_file(filepath, &json, &fsize) != 0) { + char err[256]; + snprintf(err, sizeof(err), + "{\"resp\":\"%.*s\",\"ok\":false,\"error\":\"load: %s\"}", + cmd_len, cmd, strerror(errno)); + send_line(err); + return; + } + + vlog("Loaded %ld bytes from %s (speed=%dms)\n", fsize, filepath, speed_ms); + + char resp[128]; + /* Serialize all pen access: only one stroke job touches the + * digitizer at a time, so concurrent draw/erase can't interleave. */ + pthread_mutex_lock(&g_inject_lock); + if (strncmp(cmd, "draw", 4) == 0) { + int strokes = do_draw(g_dev_fd, json, delay_us); + snprintf(resp, sizeof(resp), "{\"resp\":\"draw\",\"ok\":true,\"strokes\":%d}", strokes); + } else if (strncmp(cmd, "erase", 5) == 0) { + int points = do_erase(g_dev_fd, json, delay_us); + snprintf(resp, sizeof(resp), "{\"resp\":\"erase\",\"ok\":true,\"points\":%d}", points); + } else { + snprintf(resp, sizeof(resp), "{\"resp\":\"?\",\"ok\":false,\"error\":\"unknown cmd\"}"); + } + pthread_mutex_unlock(&g_inject_lock); + + free(json); + send_line(resp); +} + +/* ─── idle-event watcher thread ──────────────────────────────────── */ + +static void *idle_watch_thread(void *arg) { + (void)arg; + long long last_ts = -1; + + while (g_client_fd >= 0) { + FILE *fp = fopen(IDLE_PATH, "r"); + if (fp) { + long long ts = 0; + if (fscanf(fp, "%lld", &ts) == 1 && ts != last_ts) { + last_ts = ts; + char evt[64]; + snprintf(evt, sizeof(evt), "{\"event\":\"idle\",\"ts\":%lld}", ts); + send_line(evt); + vlog("Pushed idle event ts=%lld\n", ts); + } + fclose(fp); + } + usleep(200000); /* poll the local file every 200ms (no network) */ + } + return NULL; +} + +/* ─── client session ─────────────────────────────────────────────── */ + +static void serve_client(int client_fd) { + g_client_fd = client_fd; + + /* Disable Nagle so small JSON commands go out immediately */ + int one = 1; + setsockopt(client_fd, IPPROTO_TCP, TCP_NODELAY, &one, sizeof(one)); + + /* Spawn idle watcher for this session */ + pthread_t watcher; + pthread_create(&watcher, NULL, idle_watch_thread, NULL); + pthread_detach(watcher); + + send_line("{\"event\":\"ready\"}"); + vlog("Client connected\n"); + + char buf[MAX_CMD]; + char line[MAX_CMD]; + int line_len = 0; + + while (1) { + ssize_t n = read(client_fd, buf, sizeof(buf)); + if (n <= 0) break; /* client disconnected */ + + /* Accumulate into lines (newline-delimited framing) */ + for (ssize_t i = 0; i < n; i++) { + if (buf[i] == '\n') { + line[line_len] = '\0'; + if (line_len > 0) handle_command(line); + line_len = 0; + } else if (line_len < (int)sizeof(line) - 1) { + line[line_len++] = buf[i]; + } + } + } + + g_client_fd = -1; + close(client_fd); + vlog("Client disconnected\n"); +} + +/* ─── main ───────────────────────────────────────────────────────── */ + +int main(int argc, char **argv) { + int port = DEFAULT_PORT; + for (int i = 1; i < argc; i++) { + if (strcmp(argv[i], "-v") == 0 || strcmp(argv[i], "--verbose") == 0) + g_verbose = 1; + else if (strcmp(argv[i], "-p") == 0 && i + 1 < argc) + port = atoi(argv[++i]); + } + + signal(SIGPIPE, SIG_IGN); + + g_dev_fd = open(DEV_PATH, O_WRONLY); + if (g_dev_fd < 0) { + perror("open " DEV_PATH); + return 1; + } + + int srv = socket(AF_INET, SOCK_STREAM, 0); + if (srv < 0) { perror("socket"); return 1; } + + int one = 1; + setsockopt(srv, SOL_SOCKET, SO_REUSEADDR, &one, sizeof(one)); + + struct sockaddr_in addr; + memset(&addr, 0, sizeof(addr)); + addr.sin_family = AF_INET; + addr.sin_addr.s_addr = INADDR_ANY; + addr.sin_port = htons(port); + + if (bind(srv, (struct sockaddr *)&addr, sizeof(addr)) < 0) { + perror("bind"); + return 1; + } + if (listen(srv, 1) < 0) { + perror("listen"); + return 1; + } + + fprintf(stderr, "[uinjectd] Listening on :%d (dev=%s%s)\n", + port, DEV_PATH, g_verbose ? " verbose" : ""); + + /* Accept one client at a time; reconnect loop for robustness */ + while (1) { + struct sockaddr_in cli; + socklen_t clen = sizeof(cli); + int client_fd = accept(srv, (struct sockaddr *)&cli, &clen); + if (client_fd < 0) { + if (errno == EINTR) continue; + perror("accept"); + break; + } + serve_client(client_fd); + } + + close(srv); + close(g_dev_fd); + return 0; +}