From d3c914285c5a8a5dd14b6832b261042ae076c264 Mon Sep 17 00:00:00 2001 From: Daniel Dolezal Date: Mon, 27 Jul 2026 22:18:09 +0200 Subject: [PATCH] feat(cli): watch files with event streams - Add a private SSE endpoint that emits file change events for authenticated clients. - Replace the watch sleep loop with watchdog local file events and remote SSE triggers. - Debounce local event bursts and keep watch running after transient sync failures. - Add watchdog as a standalone CLI dependency and cover watch behavior in tests. - Bump NanoShare to 1.25.0 and the standalone CLI to 0.4.0. --- cli/nanoshare_client/cli.py | 9 +-- cli/nanoshare_client/client.py | 10 ++- cli/nanoshare_client/watch.py | 127 +++++++++++++++++++++++++++++++++ cli/pyproject.toml | 3 +- pyproject.toml | 2 +- routes/__init__.py | 1 + routes/api/events.py | 51 +++++++++++++ run.py | 2 + tests/test_nanoshare_cli.py | 39 ++++++++++ uv.lock | 2 +- 10 files changed, 234 insertions(+), 12 deletions(-) create mode 100644 cli/nanoshare_client/watch.py create mode 100644 routes/api/events.py diff --git a/cli/nanoshare_client/cli.py b/cli/nanoshare_client/cli.py index d3f21ca..92db59d 100644 --- a/cli/nanoshare_client/cli.py +++ b/cli/nanoshare_client/cli.py @@ -3,7 +3,6 @@ from __future__ import annotations import argparse import json import sys -import time from pathlib import Path from .app import make_client @@ -14,6 +13,7 @@ from .config import DEFAULT_CONFIG from .remote_path import join_remote_path from .sync import sync_once from .table import format_table +from .watch import watch_loop DEFAULT_NODE = 'picoshare' @@ -99,12 +99,7 @@ def _cmd_sync(args) -> int: client.close() def _cmd_watch(args) -> int: - print(f'watching {args.folder} every {args.interval}s') - while True: - code = _cmd_sync(args) - if code: - return code - time.sleep(args.interval) + return watch_loop(args, _cmd_sync) def _cmd_completion(args) -> int: try: diff --git a/cli/nanoshare_client/client.py b/cli/nanoshare_client/client.py index 81fe887..57e2d13 100644 --- a/cli/nanoshare_client/client.py +++ b/cli/nanoshare_client/client.py @@ -45,11 +45,17 @@ def list_remote(client, node: str) -> list[dict]: files = result.get('files', []) if isinstance(result, dict) else [] return files if isinstance(files, list) else [] -def download_url(client, node: str, file_id: str) -> str: +def node_url(client, node: str, path: str) -> str: base = client.registry.get(node) if not base: raise RuntimeError(f'unknown node: {node}') - return f"{base.rstrip('/')}/api/files/{quote(file_id, safe='')}/download" + return f"{base.rstrip('/')}{path}" + +def download_url(client, node: str, file_id: str) -> str: + return node_url(client, node, f"/api/files/{quote(file_id, safe='')}/download") + +def events_url(client, node: str) -> str: + return node_url(client, node, '/api/files/events') def download(client, node: str, file_id: str, path: Path) -> None: path.parent.mkdir(parents=True, exist_ok=True) diff --git a/cli/nanoshare_client/watch.py b/cli/nanoshare_client/watch.py new file mode 100644 index 0000000..c51bc1d --- /dev/null +++ b/cli/nanoshare_client/watch.py @@ -0,0 +1,127 @@ +from __future__ import annotations + +import sys +import threading +import time +from pathlib import Path +from typing import Callable + +import httpx + +from .client import events_url +from .ignore import is_ignored, load_ignore_patterns +from .sync import relative, state_path, SYNC_DIR_NAME + +SyncFunc = Callable[[object], int] + +class _ChangeFlag: + def __init__(self) -> None: + self.event = threading.Event() + + def notify(self) -> None: + self.event.set() + + def wait(self, timeout: float) -> bool: + return self.event.wait(timeout) + + def clear(self) -> None: + self.event.clear() + +def _is_relevant_path(root: Path, path: str, ignore_patterns: list[str]) -> bool: + try: + candidate = Path(path).expanduser().resolve() + rel = relative(root, candidate) + except ValueError: + return False + if candidate == state_path(root) or root / SYNC_DIR_NAME in candidate.parents: + return False + return not is_ignored(rel, ignore_patterns) + +def _remote_events_thread(args, flag: _ChangeFlag, stop: threading.Event) -> threading.Thread: + def run() -> None: + while not stop.is_set(): + client = None + try: + from .app import make_client + client = make_client(args) + headers = {'Accept': 'text/event-stream', **client.auth_headers()} + with httpx.Client(timeout=None) as http: + with http.stream('GET', events_url(client, args.node), headers=headers) as response: + response.raise_for_status() + for line in response.iter_lines(): + if stop.is_set(): + return + if line.startswith('event: files.changed'): + flag.notify() + except Exception as exc: + if not stop.is_set(): + print(f'remote event stream unavailable; retrying in 30s: {exc}', file=sys.stderr) + stop.wait(30) + finally: + if client: + client.close() + + thread = threading.Thread(target=run, daemon=True) + thread.start() + return thread + +def _watchdog_observer(root: Path, ignore_patterns: list[str], flag: _ChangeFlag): + from watchdog.events import FileSystemEventHandler + from watchdog.observers import Observer + + class Handler(FileSystemEventHandler): + def on_any_event(self, event): + paths = [event.src_path] + dest_path = getattr(event, 'dest_path', None) + if dest_path: + paths.append(dest_path) + if any(_is_relevant_path(root, path, ignore_patterns) for path in paths): + flag.notify() + + observer = Observer() + observer.schedule(Handler(), str(root), recursive=True) + observer.start() + return observer + +def watch_loop(args, sync_func: SyncFunc) -> int: + root = Path(args.folder).expanduser().resolve() + root.mkdir(parents=True, exist_ok=True) + ignore_patterns = load_ignore_patterns(root, getattr(args, 'ignore', None)) + interval = max(float(getattr(args, 'interval', 10.0)), 0.1) + flag = _ChangeFlag() + observer = None + stop_events = threading.Event() + events_thread = None + + try: + try: + observer = _watchdog_observer(root, ignore_patterns, flag) + print(f'watching {args.folder} for local changes; remote events with {interval}s fallback') + except Exception as exc: + print(f'watchdog unavailable, polling every {interval}s: {exc}', file=sys.stderr) + print(f'watching {args.folder} every {interval}s') + + events_thread = _remote_events_thread(args, flag, stop_events) + + while True: + flag.clear() + code = sync_func(args) + if code: + print(f'watch sync failed; retrying in {interval}s', file=sys.stderr) + flag.wait(interval) + continue + + if flag.wait(interval): + flag.clear() + time.sleep(1.0) + while flag.wait(1.0): + flag.clear() + except KeyboardInterrupt: + return 0 + finally: + stop_events.set() + if observer: + observer.stop() + observer.join(timeout=5) + if events_thread: + events_thread.join(timeout=5) diff --git a/cli/pyproject.toml b/cli/pyproject.toml index 7e5a391..399901b 100644 --- a/cli/pyproject.toml +++ b/cli/pyproject.toml @@ -1,11 +1,12 @@ [project] name = "nanoshare-cli" -version = "0.3.0" +version = "0.4.0" description = "NanoShare desktop CLI and folder sync client" readme = "README.md" requires-python = ">=3.11" dependencies = [ "httpx==0.28.1", + "watchdog==6.0.0", ] [project.scripts] diff --git a/pyproject.toml b/pyproject.toml index 7dbca88..d90e35f 100644 --- a/pyproject.toml +++ b/pyproject.toml @@ -1,6 +1,6 @@ [project] name = "nanoshare" -version = "1.24.0" +version = "1.25.0" description = "Add your description here" readme = "README.md" requires-python = ">=3.13" diff --git a/routes/__init__.py b/routes/__init__.py index 4988857..d7f424a 100644 --- a/routes/__init__.py +++ b/routes/__init__.py @@ -21,6 +21,7 @@ from .side.upload import upload_bp from .api.cli_auth import cli_auth_bp from .api.download import api_download_bp +from .api.events import api_events_bp # Health from .api.health import health_bp diff --git a/routes/api/events.py b/routes/api/events.py new file mode 100644 index 0000000..8297ce4 --- /dev/null +++ b/routes/api/events.py @@ -0,0 +1,51 @@ +from __future__ import annotations + +import asyncio +import json + +from my_modules.app.setup import LIMITER +from my_modules.decoratory.header import token_required +from quart import Blueprint, Response, current_app, stream_with_context + +api_events_bp = Blueprint('api_events', __name__) + +async def _snapshot(user_id: str) -> str: + files = await current_app.convex.get_files(user_id) + rows = [] + for item in files or []: + if not isinstance(item, dict): + continue + rows.append({ + 'file_id': item.get('file_id'), + 'file_name': item.get('file_name'), + 'file_path': item.get('file_path') or '', + 'file_size': item.get('file_size'), + 'expires_at': item.get('expires_at'), + 'updated_at': item.get('updated_at') or item.get('uploaded_at'), + }) + rows.sort(key=lambda row: str(row.get('file_id') or '')) + return json.dumps(rows, sort_keys=True, separators=(',', ':')) + +@api_events_bp.get('/api/files/events') +@LIMITER.limit('6 per minute;60 per hour;') +@token_required(['files', 'mesh']) +async def api_file_events(user: dict): + user_id = user['sub'] + + @stream_with_context + async def stream(): + previous = await _snapshot(user_id) + yield 'event: ready\ndata: {}\n\n' + while True: + await asyncio.sleep(30) + current = await _snapshot(user_id) + if current != previous: + previous = current + yield f'event: files.changed\ndata: {current}\n\n' + else: + yield ': keepalive\n\n' + + return Response(stream(), content_type='text/event-stream', headers={ + 'Cache-Control': 'no-cache', + 'X-Accel-Buffering': 'no', + }) diff --git a/run.py b/run.py index 181b996..2292bab 100755 --- a/run.py +++ b/run.py @@ -13,6 +13,7 @@ from routes import ( upload_bp, cli_auth_bp, api_download_bp, + api_events_bp, health_bp ) from routes.api.link import link_bp as servicelink_bp @@ -25,6 +26,7 @@ app.register_blueprint(side_main_bp) app.register_blueprint(upload_bp) app.register_blueprint(cli_auth_bp) app.register_blueprint(api_download_bp) +app.register_blueprint(api_events_bp) # ServiceLink node-to-node mesh endpoint (POST /rpc) app.register_blueprint(servicelink_bp) diff --git a/tests/test_nanoshare_cli.py b/tests/test_nanoshare_cli.py index a7ac571..094f295 100644 --- a/tests/test_nanoshare_cli.py +++ b/tests/test_nanoshare_cli.py @@ -11,6 +11,7 @@ from nanoshare_client.auth import _CallbackServer from nanoshare_client.completion import script as completion_script from nanoshare_client.ignore import load_ignore_patterns, is_ignored from nanoshare_client.table import format_table +from nanoshare_client.watch import _is_relevant_path class FakeClient: def __init__(self): @@ -268,6 +269,44 @@ def test_default_ignore_patterns_skip_common_generated_files(tmp_path): assert is_ignored('nested/Thumbs.db', patterns) assert not is_ignored('docs/note.txt', patterns) +def test_watch_loop_retries_after_sync_error(monkeypatch, tmp_path): + from nanoshare_client import watch as watch_module + + calls = [] + + class FakeFlag: + def __init__(self): + pass + + def clear(self): + pass + + def wait(self, timeout): + return False + + def fake_sync(args): + calls.append(args.folder) + if len(calls) == 1: + return 1 + raise KeyboardInterrupt + + monkeypatch.setattr(watch_module, '_ChangeFlag', FakeFlag) + monkeypatch.setattr(watch_module, '_watchdog_observer', lambda *args: None) + args = type('Args', (), {'folder': str(tmp_path), 'ignore': [], 'interval': 0.1})() + + code = watch_module.watch_loop(args, fake_sync) + + assert code == 0 + assert calls == [str(tmp_path), str(tmp_path)] + +def test_watch_relevance_ignores_sync_state_and_default_ignores(tmp_path): + patterns = load_ignore_patterns(tmp_path) + + assert _is_relevant_path(tmp_path, str(tmp_path / 'note.txt'), patterns) + assert not _is_relevant_path(tmp_path, str(tmp_path / '.nanoshare-sync' / 'state.json'), patterns) + assert not _is_relevant_path(tmp_path, str(tmp_path / '.comments' / 'comment.xml'), patterns) + assert not _is_relevant_path(tmp_path, str(tmp_path / '.env'), patterns) + def test_completion_scripts_include_commands(): zsh = completion_script('zsh') bash = completion_script('bash') diff --git a/uv.lock b/uv.lock index a163482..93c9a14 100644 --- a/uv.lock +++ b/uv.lock @@ -720,7 +720,7 @@ wheels = [ [[package]] name = "nanoshare" -version = "1.24.0" +version = "1.25.0" source = { virtual = "." } dependencies = [ { name = "aiohttp" },