From f1d497d3360e6b464490b383b44aaa2176cbe835 Mon Sep 17 00:00:00 2001 From: Hermes Agent Date: Fri, 2 Oct 2026 01:35:20 +0300 Subject: [PATCH] =?UTF-8?q?watch:=20=D1=81=D0=BB=D0=B5=D0=B6=D0=B5=D0=BD?= =?UTF-8?q?=D0=B8=D0=B5=20=D0=B7=D0=B0=20=D1=84=D0=B0=D0=B9=D0=BB=D0=B0?= =?UTF-8?q?=D0=BC=D0=B8=20(WatchService=20+=20debounce=20+=20=D1=80=D0=B5?= =?UTF-8?q?=D0=BA=D0=BE=D0=BD=D1=81=D0=B8=D0=BB=D1=8F=D1=86=D0=B8=D1=8F)?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- docs/orders/07-mcp.md | 68 +++++++++ memo-watch/build.gradle.kts | 10 ++ .../src/main/kotlin/memo/watch/WatchMain.kt | 90 +++++++++++ .../src/main/kotlin/memo/watch/Watcher.kt | 141 ++++++++++++++++++ .../src/test/kotlin/memo/watch/WatcherTest.kt | 90 +++++++++++ 5 files changed, 399 insertions(+) create mode 100644 docs/orders/07-mcp.md create mode 100644 memo-watch/src/main/kotlin/memo/watch/WatchMain.kt create mode 100644 memo-watch/src/main/kotlin/memo/watch/Watcher.kt create mode 100644 memo-watch/src/test/kotlin/memo/watch/WatcherTest.kt diff --git a/docs/orders/07-mcp.md b/docs/orders/07-mcp.md new file mode 100644 index 0000000..73eb6a6 --- /dev/null +++ b/docs/orders/07-mcp.md @@ -0,0 +1,68 @@ +Проект: /root/WORK/memo (Kotlin/JVM). Заказ по MCP-серверу. + +Создать: + memo-mcp/build.gradle.kts (kotlin("jvm") + application + зависимости на :memo-core и kotlinx-serialization-json) + memo-mcp/src/main/kotlin/memo/mcp/McpServer.kt + memo-mcp/src/test/kotlin/memo/mcp/McpProtocolTest.kt + +## Зависимости модуля memo-mcp + + plugins { kotlin("jvm") application; kotlin("plugin.serialization") version "2.4.10" } + application { mainClass.set("memo.mcp.McpServerKt") } + dependencies { + implementation(project(":memo-core")) + implementation("org.jetbrains.kotlinx:kotlinx-serialization-json:1.7.3") + } + +## McpServer.kt — протокол + +fun main(args: Array) — читает stdin построчно, отвечает в stdout ОДНОЙ строкой JSON-RPC 2.0 на запрос. +НИЧЕГО кроме JSON-RPC в stdout не печатать (диагностика — только в stderr). +Поддерживаемые методы: +- `initialize` → result: {"protocolVersion":"2024-11-05","capabilities":{"tools":{}},"serverInfo":{"name":"memo","version":"0.1.0"}} +- `notifications/initialized` → уведомление, ответ не отправлять +- `tools/list` → result: {"tools":[ ... ]} с тремя инструментами: + * memo_search: {"path":string,"query":string,"k":integer(по умолч. 8),"mode":string("hybrid"|"lex"|"vec")} + обязательные: path, query + * memo_status: {"path":string} — обязательный: path + * memo_reindex: {"path":string} — обязательный: path +- `tools/call` с params {"name":..., "arguments":{...}} → + result: {"content":[{"type":"text","text":"<строка>"}],"isError":} + Ошибка (нет обязательного аргумента, неизвестный инструмент, исключение) → result с "isError":true + и текстом ошибки, а НЕ JSON-RPC error. +- неизвестный метод → {"jsonrpc":"2.0","id":,"error":{"code":-32601,"message":"Method not found"}} + +Реализацию инструментов вынести в функции, тестируемые без процесса: + fun toolSearch(path: String, query: String, k: Int, mode: SearchMode): String + fun toolStatus(path: String): String + fun toolReindex(path: String): String + fun handleRequest(line: String): String? // null для уведомлений + +Поведение: +- Модель: MEMO_MODEL_DIR, иначе /root/WORK/memo/models/siglip2. Пути к БД — как в CLI + (<коллекция>/.memo/index.db); определение коллекций: каталог с *.md на глубине до 2. +- toolSearch: поиск через Searcher с refresh-хуком (переиндексация перед поиском). + Текст ответа — человекочитаемый список: `: \n\n` для каждого хита. + Пустой результат → строка "ничего не найдено". +- toolStatus: по коллекции — `файлов , чанков , индекс `. +- toolReindex: полная переиндексация: удалить каталог `<коллекция>/.memo` и создать заново, затем + indexTree. Вернуть `переиндексировано файлов: `. +- memo_reindex через tools/call должен РАБОТАТЬ (не isError). +- Неизвестный инструмент → isError true с текстом "unknown tool: ". + +## McpProtocolTest.kt — ровно 5 тестов, имена ровно такие + +1. `initializeHandshake` — handleRequest(строка JSON с initialize) → ответ содержит "2024-11-05" и "memo". +2. `toolsListHasThreeTools` — ответ на tools/list содержит имена memo_search, memo_status, memo_reindex. +3. `unknownMethodReturns32601` — ответ содержит "-32601". +4. `searchWithoutPathIsError` — tools/call memo_search с пустыми arguments → в ответе "isError":true. +5. `notificationProducesNoResponse` — handleRequest("{\"jsonrpc\":\"2.0\",\"method\":\"notifications/initialized\"}") == null. + +После: ./gradlew :memo-mcp:test --rerun-tasks — зелёные; ./gradlew :memo-mcp:installDist — собирается; +бинарь: memo-mcp/build/install/memo-mcp/bin/memo-mcp +Коммит: git add -A && git commit -m "mcp: сервер memo_search/memo_status/memo_reindex" + +СТРОГИЕ ЗАПРЕТЫ: +- Не выводить план текстом; сразу создавай файлы. +- Не трогать memo-core, memo-cli, memo-watch. +- Ничего не печатать в stdout, кроме строк JSON-RPC. diff --git a/memo-watch/build.gradle.kts b/memo-watch/build.gradle.kts index 9bc4506..82887d0 100644 --- a/memo-watch/build.gradle.kts +++ b/memo-watch/build.gradle.kts @@ -1,3 +1,13 @@ plugins { kotlin("jvm") + application +} + +application { + mainClass.set("memo.watch.WatchMainKt") +} + +dependencies { + implementation(project(":memo-core")) + testImplementation(kotlin("test")) } \ No newline at end of file diff --git a/memo-watch/src/main/kotlin/memo/watch/WatchMain.kt b/memo-watch/src/main/kotlin/memo/watch/WatchMain.kt new file mode 100644 index 0000000..ad1e5cb --- /dev/null +++ b/memo-watch/src/main/kotlin/memo/watch/WatchMain.kt @@ -0,0 +1,90 @@ +package memo.watch + +import memo.core.Db +import memo.core.Embedder +import memo.core.Indexer +import java.io.File +import java.util.concurrent.CountDownLatch +import kotlin.system.exitProcess + +fun main(args: Array) { + if (args.isEmpty()) { + System.err.println("usage: watch ") + exitProcess(2) + } + val raw = File(args[0]) + val base = if (raw.name == ".memo") raw.parentFile ?: raw else raw + val collections = discoverCollections(base) + if (collections.isEmpty()) { + System.err.println("коллекции не найдены в ${base.absolutePath}") + exitProcess(1) + } + + val modelDir = System.getenv("MEMO_MODEL_DIR") ?: "/root/WORK/memo/models/siglip2" + val modelPath = "$modelDir/text_model_int8.onnx" + val tokenizerPath = "$modelDir/tokenizer.model" + + val ctxList = collections.map { coll -> + val memoDir = File(coll, ".memo") + memoDir.mkdirs() + val db = Db(File(memoDir, "index.db").absolutePath) + db.init() + val embedder = Embedder(modelPath, tokenizerPath) + val indexer = Indexer(db, embedder) + CollCtx(coll, db, embedder, indexer) + } + + val indexersByCollection = ctxList.associateBy { it.root } + val watcher = Watcher( + collections = ctxList.map { it.root }, + index = { coll -> indexersByCollection.getValue(coll).indexer.indexTree(coll) }, + ) + + Runtime.getRuntime().addShutdownHook( + Thread({ + try { watcher.close() } catch (_: Throwable) {} + for (ctx in ctxList) { + try { ctx.embedder.close() } catch (_: Throwable) {} + try { ctx.db.close() } catch (_: Throwable) {} + } + }, "memo-watch-shutdown") + ) + + watcher.start() + for (ctx in ctxList) { + println("наблюдаю: ${ctx.root.absolutePath}") + } + + CountDownLatch(1).await() +} + +private data class CollCtx( + val root: File, + val db: Db, + val embedder: Embedder, + val indexer: Indexer, +) + +private fun discoverCollections(base: File): List { + if (!base.isDirectory) return emptyList() + val candidates = LinkedHashSet() + candidates.add(base) + val q = ArrayDeque>() + q.addLast(base to 0) + while (q.isNotEmpty()) { + val (d, depth) = q.removeFirst() + if (depth >= 2) continue + val children = d.listFiles() ?: continue + for (c in children) { + if (c.isDirectory && !c.name.startsWith(".")) { + candidates.add(c) + q.addLast(c to depth + 1) + } + } + } + return candidates.filter { d -> + d.walkTopDown() + .maxDepth(8) + .any { it.isFile && it.extension == "md" } + } +} \ No newline at end of file diff --git a/memo-watch/src/main/kotlin/memo/watch/Watcher.kt b/memo-watch/src/main/kotlin/memo/watch/Watcher.kt new file mode 100644 index 0000000..52e40b8 --- /dev/null +++ b/memo-watch/src/main/kotlin/memo/watch/Watcher.kt @@ -0,0 +1,141 @@ +package memo.watch + +import java.io.File +import java.nio.file.ClosedWatchServiceException +import java.nio.file.FileSystems +import java.nio.file.Path +import java.nio.file.StandardWatchEventKinds +import java.nio.file.WatchKey +import java.nio.file.WatchService +import java.util.concurrent.ConcurrentHashMap +import java.util.concurrent.TimeUnit +import java.util.concurrent.atomic.AtomicBoolean +import java.util.concurrent.atomic.AtomicInteger + +class Watcher( + private val collections: List, + private val index: (File) -> Unit, + private val debounceMillis: Long = 500, + private val reconcileMillis: Long = 10 * 60 * 1000, +) : AutoCloseable { + + private val watchService: WatchService = FileSystems.getDefault().newWatchService() + private val keyToCollection = ConcurrentHashMap() + private val running = AtomicBoolean(false) + private val closed = AtomicBoolean(false) + private val _indexCalls = AtomicInteger(0) + + val indexCalls: Int get() = _indexCalls.get() + + private val thread: Thread = Thread({ runLoop() }, "memo-watcher").apply { isDaemon = true } + + fun start() { + if (!running.compareAndSet(false, true)) return + for (coll in collections) { + try { + registerRecursive(coll, coll) + } catch (t: Throwable) { + System.err.println("watcher: register ${coll.absolutePath} failed: ${t.message}") + } + } + thread.start() + } + + private fun registerRecursive(dir: File, collection: File) { + if (!dir.isDirectory) return + if (dir.name.startsWith(".")) return + val key = dir.toPath().register( + watchService, + StandardWatchEventKinds.ENTRY_CREATE, + StandardWatchEventKinds.ENTRY_MODIFY, + StandardWatchEventKinds.ENTRY_DELETE, + ) + keyToCollection[key] = collection + val children = dir.listFiles() ?: return + for (child in children) { + if (child.isDirectory && !child.name.startsWith(".")) { + registerRecursive(child, collection) + } + } + } + + override fun close() { + if (!closed.compareAndSet(false, true)) return + running.set(false) + thread.interrupt() + try { watchService.close() } catch (_: Throwable) {} + try { thread.join(2000) } catch (_: InterruptedException) {} + } + + private fun runLoop() { + var nextReconcile = System.currentTimeMillis() + reconcileMillis + while (running.get()) { + try { + val firstKey = watchService.poll(debounceMillis, TimeUnit.MILLISECONDS) + if (firstKey != null) { + val affected = HashSet() + processKey(firstKey, affected) + val endDeadline = System.currentTimeMillis() + debounceMillis + while (running.get()) { + val now = System.currentTimeMillis() + if (now >= endDeadline) break + val remaining = endDeadline - now + val nextKey = watchService.poll(remaining, TimeUnit.MILLISECONDS) ?: break + processKey(nextKey, affected) + } + for (coll in affected) { + callIndex(coll) + } + } + if (System.currentTimeMillis() >= nextReconcile) { + nextReconcile = System.currentTimeMillis() + reconcileMillis + for (coll in collections) { + callIndex(coll) + } + } + } catch (_: ClosedWatchServiceException) { + break + } catch (_: InterruptedException) { + if (!running.get()) break + } catch (t: Throwable) { + System.err.println("watcher loop: ${t.message}") + } + } + } + + private fun processKey(key: WatchKey, affected: MutableSet) { + val collection = keyToCollection[key] ?: return + for (event in key.pollEvents()) { + val kind = event.kind() + if (kind === StandardWatchEventKinds.OVERFLOW) continue + val ctx = event.context() as? Path ?: continue + val watchable = key.watchable() as? Path ?: continue + val ev = watchable.resolve(ctx).toFile() + when (kind) { + StandardWatchEventKinds.ENTRY_CREATE -> { + affected.add(collection) + if (ev.isDirectory && !ev.name.startsWith(".")) { + try { + registerRecursive(ev, collection) + } catch (t: Throwable) { + System.err.println("watcher: registerRecursive ${ev.absolutePath} failed: ${t.message}") + } + } + } + StandardWatchEventKinds.ENTRY_MODIFY, + StandardWatchEventKinds.ENTRY_DELETE -> affected.add(collection) + } + } + val valid = key.reset() + if (!valid) keyToCollection.remove(key) + } + + private fun callIndex(coll: File) { + _indexCalls.incrementAndGet() + try { + index(coll) + } catch (t: Throwable) { + System.err.println("watcher: index ${coll.absolutePath} failed: ${t.message}") + } + } +} \ No newline at end of file diff --git a/memo-watch/src/test/kotlin/memo/watch/WatcherTest.kt b/memo-watch/src/test/kotlin/memo/watch/WatcherTest.kt new file mode 100644 index 0000000..4c0f5e5 --- /dev/null +++ b/memo-watch/src/test/kotlin/memo/watch/WatcherTest.kt @@ -0,0 +1,90 @@ +package memo.watch + +import java.io.File +import java.nio.file.Files +import java.util.concurrent.TimeUnit +import java.util.concurrent.atomic.AtomicInteger +import kotlin.test.Test +import kotlin.test.assertEquals +import kotlin.test.assertTrue + +class WatcherTest { + + @Test + fun modifyTriggersSingleIndexCall() { + val tempDir = Files.createTempDirectory("memo-watch-").toFile() + tempDir.deleteOnExit() + val note = File(tempDir, "note.md") + note.writeText("# initial") + note.deleteOnExit() + val counter = AtomicInteger(0) + val watcher = Watcher( + collections = listOf(tempDir), + index = { counter.incrementAndGet() }, + debounceMillis = 200, + reconcileMillis = TimeUnit.HOURS.toMillis(1), + ) + try { + watcher.start() + Thread.sleep(300) + note.appendText("\nappended") + assertTrue( + waitFor(5_000, 100) { counter.get() >= 1 }, + "index должен сработать в течение 5с, получено ${counter.get()}", + ) + Thread.sleep(1500) + assertEquals(1, counter.get(), "ожидался ровно 1 вызов index, получено ${counter.get()}") + } finally { + watcher.close() + tempDir.deleteRecursively() + } + } + + @Test + fun unchangedFileDoesNotTriggerIndex() { + val tempDir = Files.createTempDirectory("memo-watch-").toFile() + tempDir.deleteOnExit() + val note = File(tempDir, "note.md") + note.writeText("# initial") + note.deleteOnExit() + val counter = AtomicInteger(0) + val watcher = Watcher( + collections = listOf(tempDir), + index = { counter.incrementAndGet() }, + debounceMillis = 200, + reconcileMillis = TimeUnit.HOURS.toMillis(1), + ) + try { + watcher.start() + Thread.sleep(300) + Thread.sleep(1500) + assertEquals(0, counter.get(), "не должно быть вызовов index, получено ${counter.get()}") + } finally { + watcher.close() + tempDir.deleteRecursively() + } + } + + @Test + fun closeIsIdempotentAndStopsThread() { + val tempDir = Files.createTempDirectory("memo-watch-").toFile() + tempDir.deleteOnExit() + val watcher = Watcher( + collections = listOf(tempDir), + index = {}, + ) + watcher.start() + watcher.close() + watcher.close() + tempDir.deleteRecursively() + } + + private fun waitFor(timeoutMs: Long, stepMs: Long, condition: () -> Boolean): Boolean { + val deadline = System.currentTimeMillis() + timeoutMs + while (System.currentTimeMillis() < deadline) { + if (condition()) return true + Thread.sleep(stepMs) + } + return condition() + } +} \ No newline at end of file