diff --git a/Cargo.lock b/Cargo.lock index 7b3acff..fd6afc0 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -115,6 +115,7 @@ source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "31b698c5f9a010f6573133b09e0de5408834d0c82f8d7475a89fc1867a71cd90" dependencies = [ "axum-core", + "base64", "bytes", "form_urlencoded", "futures-util", @@ -133,8 +134,10 @@ dependencies = [ "serde_json", "serde_path_to_error", "serde_urlencoded", + "sha1", "sync_wrapper", "tokio", + "tokio-tungstenite 0.29.0", "tower", "tower-layer", "tower-service", @@ -175,6 +178,12 @@ dependencies = [ "serde", ] +[[package]] +name = "bitflags" +version = "1.3.2" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "bef38d45163c2f1dde094a7dfd33ccf595c92905c8f8f4fdc18d06fb1037718a" + [[package]] name = "bitflags" version = "2.13.1" @@ -190,6 +199,15 @@ dependencies = [ "crunchy", ] +[[package]] +name = "block-buffer" +version = "0.10.4" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "3078c7629b62d3f0439517fa394996acacc5cbc91c5a20d8c658e77abd503a71" +dependencies = [ + "generic-array", +] + [[package]] name = "bon" version = "3.10.1" @@ -268,8 +286,8 @@ source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "65c35e4b699c7e15ccbe7ee35c005e4fc0a278d22238a2857e6ce2dadeda1b06" dependencies = [ "cfg-if", - "cpufeatures", - "rand_core", + "cpufeatures 0.3.1", + "rand_core 0.10.1", ] [[package]] @@ -337,6 +355,15 @@ version = "0.8.7" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "773648b94d0e5d620f64f280777445740e61fe701025087ec8b57f45c791888b" +[[package]] +name = "cpufeatures" +version = "0.2.17" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "59ed5838eebb26a2bb2e58f6d5b5316989ae9d08bab10e0e6d103e656d1b0280" +dependencies = [ + "libc", +] + [[package]] name = "cpufeatures" version = "0.3.1" @@ -395,6 +422,16 @@ version = "0.2.4" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "460fbee9c2c2f33933d720630a6a0bac33ba7053db5344fac858d4b8952d77d5" +[[package]] +name = "crypto-common" +version = "0.1.7" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "78c8292055d1c1df0cce5d180393dc8cce0abec0a7102adb6c7b1eef6016d60a" +dependencies = [ + "generic-array", + "typenum", +] + [[package]] name = "darling" version = "0.24.1" @@ -443,6 +480,12 @@ dependencies = [ "parking_lot_core", ] +[[package]] +name = "data-encoding" +version = "2.11.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "4583a4551df46e2792f82ceeac45e850d2e2d5debba0b91f102385cda5b11f06" + [[package]] name = "datasketches" version = "0.2.0" @@ -458,6 +501,16 @@ dependencies = [ "serde_core", ] +[[package]] +name = "digest" +version = "0.10.7" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "9ed9a281f7bc9b7576e61468ba615a66a5c8cfdff42420a70aa82701a3b1e292" +dependencies = [ + "block-buffer", + "crypto-common", +] + [[package]] name = "dirs" version = "6.0.0" @@ -547,6 +600,16 @@ version = "2.5.0" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "da7c62ceae207dd37ea5b845da6a0696c799f85e97da1ab5b7910be3c1c80223" +[[package]] +name = "filetime" +version = "0.2.29" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "5c287a33c7f0a620c38e641e7f60827713987b3c0f26e8ddc9462cc69cf75759" +dependencies = [ + "cfg-if", + "libc", +] + [[package]] name = "find-msvc-tools" version = "0.1.12" @@ -584,6 +647,15 @@ dependencies = [ "windows-sys 0.59.0", ] +[[package]] +name = "fsevent-sys" +version = "4.1.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "76ee7a02da4d231650c7cea31349b889be2f45ddb3ef3032d2ec8185f6313fd2" +dependencies = [ + "libc", +] + [[package]] name = "futures-channel" version = "0.3.34" @@ -644,6 +716,16 @@ dependencies = [ "slab", ] +[[package]] +name = "generic-array" +version = "0.14.7" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "85649ca51fd72272d7821adaf274ad91c288277713d9c18820d8499a7ff69e9a" +dependencies = [ + "typenum", + "version_check", +] + [[package]] name = "getrandom" version = "0.2.17" @@ -657,6 +739,18 @@ dependencies = [ "wasm-bindgen", ] +[[package]] +name = "getrandom" +version = "0.3.4" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "899def5c37c4fd7b2664648c28120ecec138e4d395b459e5ca34f9cce2dd77fd" +dependencies = [ + "cfg-if", + "libc", + "r-efi 5.3.0", + "wasip2", +] + [[package]] name = "getrandom" version = "0.4.3" @@ -666,11 +760,24 @@ dependencies = [ "cfg-if", "js-sys", "libc", - "r-efi", - "rand_core", + "r-efi 6.0.0", + "rand_core 0.10.1", "wasm-bindgen", ] +[[package]] +name = "git2" +version = "0.19.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "b903b73e45dc0c6c596f2d37eccece7c1c8bb6e4407b001096387c63d0d93724" +dependencies = [ + "bitflags 2.13.1", + "libc", + "libgit2-sys", + "log", + "url", +] + [[package]] name = "glob" version = "0.3.4" @@ -945,6 +1052,26 @@ dependencies = [ "icu_properties", ] +[[package]] +name = "inotify" +version = "0.9.6" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "f8069d3ec154eb856955c1c0fbffefbf5f3c40a104ec912d4797314c1801abff" +dependencies = [ + "bitflags 1.3.2", + "inotify-sys", + "libc", +] + +[[package]] +name = "inotify-sys" +version = "0.1.8" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "c033f80b2c113cdf91ab7a33faa9cbc014726dcad99880c8609af2a370edf37d" +dependencies = [ + "libc", +] + [[package]] name = "inventory" version = "0.3.24" @@ -1002,6 +1129,26 @@ dependencies = [ "wasm-bindgen", ] +[[package]] +name = "kqueue" +version = "1.2.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "8d763e5b24120b4ddf50de6c92308156765aabfbbccebf401da7cff2d70a41ea" +dependencies = [ + "kqueue-sys", + "libc", +] + +[[package]] +name = "kqueue-sys" +version = "1.1.2" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "07293a4e297ac234359b510362495713f75ea345d5307140414f20c69ffeb087" +dependencies = [ + "bitflags 2.13.1", + "libc", +] + [[package]] name = "lazy_static" version = "1.5.0" @@ -1020,6 +1167,18 @@ version = "0.2.189" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "3eaf3ede3fee6db1a4c2ee091bf8a8b4dccdc6d17f656fb07896ee72867612f2" +[[package]] +name = "libgit2-sys" +version = "0.17.0+1.8.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "10472326a8a6477c3c20a64547b0059e4b0d086869eee31e6d7da728a8eb7224" +dependencies = [ + "cc", + "libc", + "libz-sys", + "pkg-config", +] + [[package]] name = "libredox" version = "0.1.23" @@ -1029,6 +1188,18 @@ dependencies = [ "libc", ] +[[package]] +name = "libz-sys" +version = "1.1.29" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "85bc9657773828b90eeb625adff10eeac83cc21bbfd8e23a03eaa8a33c9e28d9" +dependencies = [ + "cc", + "libc", + "pkg-config", + "vcpkg", +] + [[package]] name = "linux-raw-sys" version = "0.12.1" @@ -1107,7 +1278,9 @@ dependencies = [ "dashmap", "dirs", "futures-util", + "git2", "glob", + "notify", "redb", "reqwest", "schemars 1.2.2", @@ -1116,6 +1289,7 @@ dependencies = [ "tantivy", "tokio", "tokio-stream", + "tokio-tungstenite 0.21.0", "tokio-util", "tracing", "tracing-subscriber", @@ -1130,6 +1304,7 @@ dependencies = [ "futures-util", "reqwest", "tokio", + "tokio-tungstenite 0.21.0", "tokio-util", ] @@ -1181,6 +1356,18 @@ version = "0.2.1" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "68354c5c6bd36d73ff3feceb05efa59b6acb7626617f4962be322a825e61f79a" +[[package]] +name = "mio" +version = "0.8.11" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "a4a650543ca06a924e8b371db273b2756685faae30f8487da1b56505a8f78b0c" +dependencies = [ + "libc", + "log", + "wasi", + "windows-sys 0.48.0", +] + [[package]] name = "mio" version = "1.2.3" @@ -1208,6 +1395,25 @@ dependencies = [ "minimal-lexical", ] +[[package]] +name = "notify" +version = "6.1.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "6205bd8bb1e454ad2e27422015fb5e4f2bcc7e08fa8f27058670d208324a4d2d" +dependencies = [ + "bitflags 2.13.1", + "crossbeam-channel", + "filetime", + "fsevent-sys", + "inotify", + "kqueue", + "libc", + "log", + "mio 0.8.11", + "walkdir", + "windows-sys 0.48.0", +] + [[package]] name = "nu-ansi-term" version = "0.50.3" @@ -1330,6 +1536,15 @@ version = "0.2.0" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "439ee305def115ba05938db6eb1644ff94165c5ab5e9420d1c1bcedbba909391" +[[package]] +name = "ppv-lite86" +version = "0.2.21" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "85eae3c4ed2f50dcfe72643da4befc30deadb458a9b590d720cde2f2b1e97da9" +dependencies = [ + "zerocopy", +] + [[package]] name = "prettyplease" version = "0.3.0" @@ -1363,7 +1578,7 @@ dependencies = [ "rustc-hash", "rustls", "socket2", - "thiserror", + "thiserror 2.0.20", "tokio", "tracing", "web-time", @@ -1378,14 +1593,14 @@ dependencies = [ "bytes", "getrandom 0.4.3", "lru-slab", - "rand", + "rand 0.10.2", "rand_pcg", "ring", "rustc-hash", "rustls", "rustls-pki-types", "slab", - "thiserror", + "thiserror 2.0.20", "tinyvec", "tracing", "web-time", @@ -1414,12 +1629,39 @@ dependencies = [ "proc-macro2", ] +[[package]] +name = "r-efi" +version = "5.3.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "69cdb34c158ceb288df11e18b4bd39de994f6657d83847bdffdbd7f346754b0f" + [[package]] name = "r-efi" version = "6.0.0" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "f8dcc9c7d52a811697d2151c701e0d08956f92b0e24136cf4cf27b57a6a0d9bf" +[[package]] +name = "rand" +version = "0.8.8" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "e058c7de0b26af77780c769414d6257830bb240f3c38477dbc2c16e5f54d6d4c" +dependencies = [ + "libc", + "rand_chacha 0.3.1", + "rand_core 0.6.4", +] + +[[package]] +name = "rand" +version = "0.9.5" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "b9ef1d0d795eb7d84685bca4f72f3649f064e6641543d3a8c415898726a57b41" +dependencies = [ + "rand_chacha 0.9.0", + "rand_core 0.9.5", +] + [[package]] name = "rand" version = "0.10.2" @@ -1428,7 +1670,45 @@ checksum = "c7f5fa3a058cd35567ef9bfa5e75732bee0f9e4c55fa90477bef2dfcdbc4be80" dependencies = [ "chacha20", "getrandom 0.4.3", - "rand_core", + "rand_core 0.10.1", +] + +[[package]] +name = "rand_chacha" +version = "0.3.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "e6c10a63a0fa32252be49d21e7709d4d4baf8d231c2dbce1eaa8141b9b127d88" +dependencies = [ + "ppv-lite86", + "rand_core 0.6.4", +] + +[[package]] +name = "rand_chacha" +version = "0.9.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "d3022b5f1df60f26e1ffddd6c66e8aa15de382ae63b3a0c1bfc0e4d3e3f325cb" +dependencies = [ + "ppv-lite86", + "rand_core 0.9.5", +] + +[[package]] +name = "rand_core" +version = "0.6.4" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "ec0be4795e2f6a28069bec0b5ff3e2ac9bafc99e6a9a7dc3547996c5c816922c" +dependencies = [ + "getrandom 0.2.17", +] + +[[package]] +name = "rand_core" +version = "0.9.5" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "76afc826de14238e6e8c374ddcc1fa19e374fd8dd986b0d2af0d02377261d83c" +dependencies = [ + "getrandom 0.3.4", ] [[package]] @@ -1443,7 +1723,7 @@ version = "0.10.2" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "caa0f4137e1c0a72f4c651489402276c8e8e1cf081f3b0ba156d2cbeef09e86a" dependencies = [ - "rand_core", + "rand_core 0.10.1", ] [[package]] @@ -1481,7 +1761,7 @@ version = "0.5.18" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "ed2bf2547551a7053d6fdfafda3f938979645c44812fbfcda098faae3f1a362d" dependencies = [ - "bitflags", + "bitflags 2.13.1", ] [[package]] @@ -1492,7 +1772,7 @@ checksum = "a4e608c6638b9c18977b00b475ac1f28d14e84b27d8d42f70e0bf1e3dec127ac" dependencies = [ "getrandom 0.2.17", "libredox", - "thiserror", + "thiserror 2.0.20", ] [[package]] @@ -1649,7 +1929,7 @@ version = "1.1.4" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "b6fe4565b9518b83ef4f91bb47ce29620ca828bd32cb7e408f0062e9930ba190" dependencies = [ - "bitflags", + "bitflags 2.13.1", "errno", "libc", "linux-raw-sys", @@ -1703,6 +1983,15 @@ version = "1.0.23" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "9774ba4a74de5f7b1c1451ed6cd5285a32eddb5cccb8cc655a4e50009e06477f" +[[package]] +name = "same-file" +version = "1.0.6" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "93fc1dc3aaa9bfed95e02e6eadabb4baf7e3078b0bd1b4d7b6b0b68378900502" +dependencies = [ + "winapi-util", +] + [[package]] name = "schemars" version = "0.8.22" @@ -1846,6 +2135,17 @@ dependencies = [ "serde", ] +[[package]] +name = "sha1" +version = "0.10.7" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "a978451301f4db1d02937a4ab3ccce137717b81826e79b7d49ffe3244a13c3b8" +dependencies = [ + "cfg-if", + "cpufeatures 0.2.17", + "digest", +] + [[package]] name = "sharded-slab" version = "0.1.7" @@ -2008,7 +2308,7 @@ dependencies = [ "tantivy-stacker", "tantivy-tokenizer-api", "tempfile", - "thiserror", + "thiserror 2.0.20", "time", "typetag", "uuid", @@ -2123,13 +2423,33 @@ dependencies = [ "windows-sys 0.61.2", ] +[[package]] +name = "thiserror" +version = "1.0.69" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "b6aaf5339b578ea85b50e080feb250a3e8ae8cfcdff9a461c9ec2904bc923f52" +dependencies = [ + "thiserror-impl 1.0.69", +] + [[package]] name = "thiserror" version = "2.0.20" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "ec86235f5fcc2a73650310756d2ac5b138a5780bbbdfae3eeccec992c435ba4f" dependencies = [ - "thiserror-impl", + "thiserror-impl 2.0.20", +] + +[[package]] +name = "thiserror-impl" +version = "1.0.69" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "4fee6c4efc90059e10f81e6d42c60a18f76588c3d74cb83a0b242a2b6c7504c1" +dependencies = [ + "proc-macro2", + "quote", + "syn 2.0.119", ] [[package]] @@ -2215,7 +2535,7 @@ checksum = "202caea871b69668250d242070849eb495be178ed697a3e98aebce5bc81a0bed" dependencies = [ "bytes", "libc", - "mio", + "mio 1.2.3", "parking_lot", "pin-project-lite", "signal-hook-registry", @@ -2256,6 +2576,30 @@ dependencies = [ "tokio", ] +[[package]] +name = "tokio-tungstenite" +version = "0.21.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "c83b561d025642014097b66e6c1bb422783339e0909e4429cde4749d1990bc38" +dependencies = [ + "futures-util", + "log", + "tokio", + "tungstenite 0.21.0", +] + +[[package]] +name = "tokio-tungstenite" +version = "0.29.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "8f72a05e828585856dacd553fba484c242c46e391fb0e58917c942ee9202915c" +dependencies = [ + "futures-util", + "log", + "tokio", + "tungstenite 0.29.0", +] + [[package]] name = "tokio-util" version = "0.7.19" @@ -2291,7 +2635,7 @@ version = "0.6.11" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "4cfcf7e2740e6fc6d4d688b4ef00650406bb94adf4731e43c096c3a19fe40840" dependencies = [ - "bitflags", + "bitflags 2.13.1", "bytes", "futures-util", "http", @@ -2379,12 +2723,53 @@ version = "0.2.5" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "e421abadd41a4225275504ea4d6566923418b7f05506fbc9c0fe86ba7396114b" +[[package]] +name = "tungstenite" +version = "0.21.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "9ef1a641ea34f399a848dea702823bbecfb4c486f911735368f1f137cb8257e1" +dependencies = [ + "byteorder", + "bytes", + "data-encoding", + "http", + "httparse", + "log", + "rand 0.8.8", + "sha1", + "thiserror 1.0.69", + "url", + "utf-8", +] + +[[package]] +name = "tungstenite" +version = "0.29.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "6c01152af293afb9c7c2a57e4b559c5620b421f6d133261c60dd2d0cdb38e6b8" +dependencies = [ + "bytes", + "data-encoding", + "http", + "httparse", + "log", + "rand 0.9.5", + "sha1", + "thiserror 2.0.20", +] + [[package]] name = "typeid" version = "1.0.3" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "bc7d623258602320d5c55d1bc22793b57daff0ec7efc270ea7d55ce1d5f5471c" +[[package]] +name = "typenum" +version = "1.20.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "b6f5e870be6c3b371b77fe0ee0bafb859fa4964b4404c27de1d380043c4dda20" + [[package]] name = "typetag" version = "0.2.23" @@ -2433,6 +2818,12 @@ dependencies = [ "serde", ] +[[package]] +name = "utf-8" +version = "0.7.6" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "09cc8ee72d2a9becf2f2febe0205bbed8fc6615b7cb429ad062dc7b7ddd036a9" + [[package]] name = "utf8-ranges" version = "1.0.5" @@ -2469,6 +2860,28 @@ version = "0.1.1" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "ba73ea9cf16a25df0c8caa16c51acb937d5712a8429db78a3ee29d5dcacd3a65" +[[package]] +name = "vcpkg" +version = "0.2.15" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "accd4ea62f7bb7a82fe23066fb0957d48ef677f6eeb8215f372f52e48bb32426" + +[[package]] +name = "version_check" +version = "0.9.5" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "0b928f33d975fc6ad9f86c8f283853ad26bdd5b10b7f1542aa2fa15e2289105a" + +[[package]] +name = "walkdir" +version = "2.5.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "29790946404f91d9c5d06f9874efddea1dc06c5efe94541a7d6863108e3a5e4b" +dependencies = [ + "same-file", + "winapi-util", +] + [[package]] name = "want" version = "0.3.1" @@ -2484,6 +2897,15 @@ version = "0.11.1+wasi-snapshot-preview1" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "ccf3ec651a847eb01de73ccad15eb7d99f80485de043efb2f370cd654f4ea44b" +[[package]] +name = "wasip2" +version = "1.0.4+wasi-0.2.12" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "b67efb37e106e55ce722a510d6b5f9c17f083e5fc79afc2badeb12cc313d9487" +dependencies = [ + "wit-bindgen", +] + [[package]] name = "wasm-bindgen" version = "0.2.128" @@ -2597,6 +3019,15 @@ version = "0.4.0" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "ac3b87c63620426dd9b991e5ce0329eff545bccbbb34f3be09ff6fb6ab51b7b6" +[[package]] +name = "winapi-util" +version = "0.1.11" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "c2a7b1c03c876122aa43f3020e6c3c3ee5c05081c9a00739faf7503aeba10d22" +dependencies = [ + "windows-sys 0.61.2", +] + [[package]] name = "winapi-x86_64-pc-windows-gnu" version = "0.4.0" @@ -2662,13 +3093,22 @@ dependencies = [ "windows-link", ] +[[package]] +name = "windows-sys" +version = "0.48.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "677d2418bec65e3338edb076e806bc1ec15693c5d0104683f2efe857f61056a9" +dependencies = [ + "windows-targets 0.48.5", +] + [[package]] name = "windows-sys" version = "0.52.0" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "282be5f36a8ce781fad8c8ae18fa3f9beff57ec1b52cb3de0789201425d9a33d" dependencies = [ - "windows-targets", + "windows-targets 0.52.6", ] [[package]] @@ -2677,7 +3117,7 @@ version = "0.59.0" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "1e38bc4d79ed67fd075bcc251a1c39b32a1776bbe92e5bef1f0bf1f8c531853b" dependencies = [ - "windows-targets", + "windows-targets 0.52.6", ] [[package]] @@ -2689,34 +3129,67 @@ dependencies = [ "windows-link", ] +[[package]] +name = "windows-targets" +version = "0.48.5" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "9a2fa6e2155d7247be68c096456083145c183cbbbc2764150dda45a87197940c" +dependencies = [ + "windows_aarch64_gnullvm 0.48.5", + "windows_aarch64_msvc 0.48.5", + "windows_i686_gnu 0.48.5", + "windows_i686_msvc 0.48.5", + "windows_x86_64_gnu 0.48.5", + "windows_x86_64_gnullvm 0.48.5", + "windows_x86_64_msvc 0.48.5", +] + [[package]] name = "windows-targets" version = "0.52.6" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "9b724f72796e036ab90c1021d4780d4d3d648aca59e491e6b98e725b84e99973" dependencies = [ - "windows_aarch64_gnullvm", - "windows_aarch64_msvc", - "windows_i686_gnu", + "windows_aarch64_gnullvm 0.52.6", + "windows_aarch64_msvc 0.52.6", + "windows_i686_gnu 0.52.6", "windows_i686_gnullvm", - "windows_i686_msvc", - "windows_x86_64_gnu", - "windows_x86_64_gnullvm", - "windows_x86_64_msvc", + "windows_i686_msvc 0.52.6", + "windows_x86_64_gnu 0.52.6", + "windows_x86_64_gnullvm 0.52.6", + "windows_x86_64_msvc 0.52.6", ] +[[package]] +name = "windows_aarch64_gnullvm" +version = "0.48.5" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "2b38e32f0abccf9987a4e3079dfb67dcd799fb61361e53e2882c3cbaf0d905d8" + [[package]] name = "windows_aarch64_gnullvm" version = "0.52.6" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "32a4622180e7a0ec044bb555404c800bc9fd9ec262ec147edd5989ccd0c02cd3" +[[package]] +name = "windows_aarch64_msvc" +version = "0.48.5" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "dc35310971f3b2dbbf3f0690a219f40e2d9afcf64f9ab7cc1be722937c26b4bc" + [[package]] name = "windows_aarch64_msvc" version = "0.52.6" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "09ec2a7bb152e2252b53fa7803150007879548bc709c039df7627cabbd05d469" +[[package]] +name = "windows_i686_gnu" +version = "0.48.5" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "a75915e7def60c94dcef72200b9a8e58e5091744960da64ec734a6c6e9b3743e" + [[package]] name = "windows_i686_gnu" version = "0.52.6" @@ -2729,30 +3202,60 @@ version = "0.52.6" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "0eee52d38c090b3caa76c563b86c3a4bd71ef1a819287c19d586d7334ae8ed66" +[[package]] +name = "windows_i686_msvc" +version = "0.48.5" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "8f55c233f70c4b27f66c523580f78f1004e8b5a8b659e05a4eb49d4166cca406" + [[package]] name = "windows_i686_msvc" version = "0.52.6" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "240948bc05c5e7c6dabba28bf89d89ffce3e303022809e73deaefe4f6ec56c66" +[[package]] +name = "windows_x86_64_gnu" +version = "0.48.5" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "53d40abd2583d23e4718fddf1ebec84dbff8381c07cae67ff7768bbf19c6718e" + [[package]] name = "windows_x86_64_gnu" version = "0.52.6" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "147a5c80aabfbf0c7d901cb5895d1de30ef2907eb21fbbab29ca94c5b08b1a78" +[[package]] +name = "windows_x86_64_gnullvm" +version = "0.48.5" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "0b7b52767868a23d5bab768e390dc5f5c55825b6d30b86c844ff2dc7414044cc" + [[package]] name = "windows_x86_64_gnullvm" version = "0.52.6" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "24d5b23dc417412679681396f2b49f3de8c1473deb516bd34410872eff51ed0d" +[[package]] +name = "windows_x86_64_msvc" +version = "0.48.5" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "ed94fce61571a4006852b7389a063ab983c02eb1bb37b47f8272ce92d06d9538" + [[package]] name = "windows_x86_64_msvc" version = "0.52.6" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "589f6da84c646204747d1270a2a5661ea66ed1cced2631d546fdfb155959f9ec" +[[package]] +name = "wit-bindgen" +version = "0.57.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "1ebf944e87a7c253233ad6766e082e3cd714b5d03812acc24c318f549614536e" + [[package]] name = "writeable" version = "0.6.4" @@ -2782,6 +3285,26 @@ dependencies = [ "synstructure", ] +[[package]] +name = "zerocopy" +version = "0.8.57" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "d35102a9f36d089ccae9e4c6802bc118be4487b80aaffc0ab4e0cf5ce92d2873" +dependencies = [ + "zerocopy-derive", +] + +[[package]] +name = "zerocopy-derive" +version = "0.8.57" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "146c01f5ab44258da43cf276c74a2763db2ff3969c9c652c3f2de07041d0b2bc" +dependencies = [ + "proc-macro2", + "quote", + "syn 2.0.119", +] + [[package]] name = "zerofrom" version = "0.1.8" diff --git a/fix_all.py b/fix_all.py new file mode 100644 index 0000000..aa072ac --- /dev/null +++ b/fix_all.py @@ -0,0 +1,25 @@ +import sys + +handlers_path = r"C:\Users\reazul.ashraf\workspace\rust\mcp-memory\server\src\handlers.rs" +with open(handlers_path, "r", encoding="utf-8") as f: + content = f.read() + +content = content.replace( + 'return Some(crate::mcp::success(id.clone(), serde_json::json!({"isError": true, "content": [{"type": "text", "text": format!("Invalid args: {}", e)}] })));', + 'return Some(crate::mcp::success(id.clone(), serde_json::json!({"isError": true, "content": [{"type": "text", "text": format!("Invalid args: {}", e)}] })))' +) +content = content.replace('let graph = self.state.graph.read();', 'let graph = self.state.get_full_graph();') + +with open(handlers_path, "w", encoding="utf-8") as f: + f.write(content) +print("Fixed handlers.rs again") + +main_path = r"C:\Users\reazul.ashraf\workspace\rust\mcp-memory\server\src\main.rs" +with open(main_path, "r", encoding="utf-8") as f: + main_content = f.read() + +main_content = main_content.replace('state_git.code_changes.modify(|changes| {', 'state_git.ledger.modify(|changes| {') + +with open(main_path, "w", encoding="utf-8") as f: + f.write(main_content) +print("Fixed main.rs again") diff --git a/fix_bg.py b/fix_bg.py new file mode 100644 index 0000000..77e5dd6 --- /dev/null +++ b/fix_bg.py @@ -0,0 +1,186 @@ +import re + +with open("server/src/main.rs", "r", encoding="utf-8") as f: + content = f.read() + +# 1. Remove ledger and sticky GC from reconcile_worker +reconcile_worker_old = """ + let now = SystemTime::now() + .duration_since(UNIX_EPOCH) + .unwrap() + .as_secs(); + state.ledger.modify(|ledger| { + let seven_days = now.saturating_sub(7 * 24 * 60 * 60); + ledger.retain(|c| c.timestamp >= seven_days); + if ledger.len() > 1000 { + let excess = ledger.len() - 1000; + ledger.drain(0..excess); + } + }); + state.sticky.modify(|notes| { + notes.retain(|note| note.timestamp >= now.saturating_sub(24 * 60 * 60)); + });""" + +content = content.replace(reconcile_worker_old, "") + +# 2. Extract Task GC from main and create garbage_collector_worker +gc_loop_old = """ // Background Garbage Collection for old tasks + let state_gc = Arc::clone(&state); + tokio::spawn(async move { + loop { + // Run every 24 hours + tokio::time::sleep(tokio::time::Duration::from_secs(24 * 3600)).await; + + let now = std::time::SystemTime::now().duration_since(std::time::UNIX_EPOCH).unwrap().as_secs(); + let fourteen_days = 14 * 24 * 3600; + let cutoff = now.saturating_sub(fourteen_days); + + state_gc.tasks.modify(|tasks| { + let initial_len = tasks.len(); + tasks.retain(|task| { + if task.status.to_lowercase() == "completed" && task.created_at < cutoff { + false // remove + } else { + true // keep + } + }); + if tasks.len() < initial_len { + eprintln!("GC: Removed {} old completed tasks", initial_len - tasks.len()); + } + }); + } + });""" + +gc_worker_new = """async fn garbage_collector_worker(state: Arc) { + loop { + // Run every 6 hours + tokio::time::sleep(tokio::time::Duration::from_secs(6 * 3600)).await; + + let now = std::time::SystemTime::now().duration_since(std::time::UNIX_EPOCH).unwrap().as_secs(); + + // 1. Task GC (14 days) + let fourteen_days = 14 * 24 * 3600; + let task_cutoff = now.saturating_sub(fourteen_days); + state.tasks.modify(|tasks| { + let initial_len = tasks.len(); + tasks.retain(|task| !(task.status.to_lowercase() == "completed" && task.created_at < task_cutoff)); + if tasks.len() < initial_len { + eprintln!("GC: Removed {} old completed tasks", initial_len - tasks.len()); + } + }); + + // 2. Ledger GC (7 days or max 1000 items) + state.ledger.modify(|ledger| { + let seven_days = now.saturating_sub(7 * 24 * 3600); + ledger.retain(|c| c.timestamp >= seven_days); + if ledger.len() > 1000 { + let excess = ledger.len() - 1000; + ledger.drain(0..excess); + } + }); + + // 3. Sticky Notes GC (24 hours) + state.sticky.modify(|notes| { + notes.retain(|note| note.timestamp >= now.saturating_sub(24 * 3600)); + }); + } +} +""" + +content = content.replace(gc_loop_old, " tokio::spawn(garbage_collector_worker(Arc::clone(&state)));") +content = content.replace("async fn reconcile_worker", gc_worker_new + "\nasync fn reconcile_worker") + +# 3. Extract Git Native Sync from main +git_loop_old = """ // Git Native Sync Background Task + let state_git = Arc::clone(&state); + tokio::spawn(async move { + let repo_path = std::env::current_dir().unwrap_or_else(|_| ".".into()); + let mut last_commit_id = String::new(); + + loop { + tokio::time::sleep(tokio::time::Duration::from_secs(30)).await; + + if let Ok(repo) = git2::Repository::discover(&repo_path) { + if let Ok(head) = repo.head() { + if let Ok(commit) = head.peel_to_commit() { + let current_id = commit.id().to_string(); + if current_id != last_commit_id && !last_commit_id.is_empty() { + let msg = commit.message().unwrap_or("").to_string(); + let branch = head.shorthand().unwrap_or("unknown").to_string(); + + state_git.ledger.modify(|changes| { + changes.push(crate::models::CodeChange { + git_commit: Some(current_id.clone()), + git_branch: Some(branch), + description: format!("Auto-synced commit: {}", msg.trim()), + timestamp: std::time::SystemTime::now().duration_since(std::time::UNIX_EPOCH).unwrap().as_secs(), + file_path: "".to_string(), + }); + }); + eprintln!("Git Sync: Logged new commit {}", current_id); + + state_git.tasks.modify(|tasks| { + for task in tasks.iter_mut() { + if task.status != "completed" && msg.to_lowercase().contains(&task.title.to_lowercase()) { + task.status = "completed".to_string(); + eprintln!("Git Sync: Auto-completed task '{}'", task.title); + } + } + }); + } + last_commit_id = current_id; + } + } + } + } + });""" + +git_worker_new = """async fn git_sync_worker(state: Arc) { + let repo_path = std::env::current_dir().unwrap_or_else(|_| ".".into()); + let mut last_commit_id = String::new(); + + loop { + tokio::time::sleep(tokio::time::Duration::from_secs(30)).await; + + if let Ok(repo) = git2::Repository::discover(&repo_path) { + if let Ok(head) = repo.head() { + if let Ok(commit) = head.peel_to_commit() { + let current_id = commit.id().to_string(); + if current_id != last_commit_id && !last_commit_id.is_empty() { + let msg = commit.message().unwrap_or("").to_string(); + let branch = head.shorthand().unwrap_or("unknown").to_string(); + + state.ledger.modify(|changes| { + changes.push(crate::models::CodeChange { + git_commit: Some(current_id.clone()), + git_branch: Some(branch), + description: format!("Auto-synced commit: {}", msg.trim()), + timestamp: std::time::SystemTime::now().duration_since(std::time::UNIX_EPOCH).unwrap().as_secs(), + file_path: "".to_string(), + }); + }); + eprintln!("Git Sync: Logged new commit {}", current_id); + + state.tasks.modify(|tasks| { + for task in tasks.iter_mut() { + if task.status != "completed" && msg.to_lowercase().contains(&task.title.to_lowercase()) { + task.status = "completed".to_string(); + eprintln!("Git Sync: Auto-completed task '{}'", task.title); + } + } + }); + } + last_commit_id = current_id; + } + } + } + } +} +""" + +content = content.replace(git_loop_old, " tokio::spawn(git_sync_worker(Arc::clone(&state)));") +content = content.replace("async fn reconcile_worker", git_worker_new + "\nasync fn reconcile_worker") + +with open("server/src/main.rs", "w", encoding="utf-8") as f: + f.write(content) + diff --git a/fix_cli.py b/fix_cli.py new file mode 100644 index 0000000..cf54569 --- /dev/null +++ b/fix_cli.py @@ -0,0 +1,97 @@ +import re + +with open("server/src/main.rs", "r", encoding="utf-8") as f: + content = f.read() + +cli_gate_block = re.search(r' if let Some\(command\) = cli\.command \{.*?(?= run_server\(state\))', content, re.DOTALL) +if cli_gate_block: + content = content.replace(cli_gate_block.group(0), "") + +new_cli_logic = """ + if let Some(command) = cli.command { + let rt = tokio::runtime::Runtime::new().unwrap(); + rt.block_on(async { + let client = reqwest::Client::new(); + match command { + Commands::Gate { subcmd } => match subcmd { + GateCommands::Set { action, target, namespace, params, authorize, block, reason } => { + let mut param_map = HashMap::new(); + for p in params { + if let Some((k, v)) = p.split_once('=') { + param_map.insert(k.to_string(), v.to_string()); + } + } + let body = serde_json::json!({ + "action": action, + "target": target, + "namespace": namespace, + "params": param_map, + "authorize": authorize, + "block": block, + "reason": reason + }); + match client.post("http://127.0.0.1:3000/gate/set").json(&body).send().await { + Ok(res) if res.status().is_success() => { + println!("Gate state updated via daemon."); + std::process::exit(0); + } + Ok(res) => { + eprintln!("Failed to update gate: {}", res.status()); + std::process::exit(1); + } + Err(e) => { + eprintln!("Failed to connect to daemon: {}", e); + std::process::exit(1); + } + } + } + GateCommands::Verify { action, target, namespace, params, consume } => { + let mut query = vec![ + ("action".to_string(), action), + ("target".to_string(), target), + ("consume".to_string(), consume.to_string()), + ]; + if let Some(ns) = namespace { + query.push(("namespace".to_string(), ns)); + } + // Reqwest will serialize params correctly if we pass the right struct. + // Actually, axum's Query extractor for HashMap requires flat keys or standard serialization. + for p in params { + if let Some((k, v)) = p.split_once('=') { + query.push((k.to_string(), v.to_string())); + } + } + match client.get("http://127.0.0.1:3000/gate/verify").query(&query).send().await { + Ok(res) => { + let status = res.status(); + let text = res.text().await.unwrap_or_default(); + if status.is_success() { + std::process::exit(0); + } else if status == reqwest::StatusCode::FORBIDDEN { + eprintln!("❌ {}", text); + std::process::exit(1); + } else { + eprintln!("❌ {}", text); + std::process::exit(2); + } + } + Err(e) => { + eprintln!("Failed to connect to daemon: {}", e); + std::process::exit(2); + } + } + } + } + } + }); + return Ok(()); + } +""" + +restart_block = re.search(r' if cli\.restart \{.*?return Ok\(\(\);\n \}', content, re.DOTALL) +if restart_block: + content = content.replace(restart_block.group(0), restart_block.group(0) + "\n" + new_cli_logic) + +with open("server/src/main.rs", "w", encoding="utf-8") as f: + f.write(content) + diff --git a/fix_git.py b/fix_git.py new file mode 100644 index 0000000..609c8ad --- /dev/null +++ b/fix_git.py @@ -0,0 +1,102 @@ +import re + +with open("server/src/main.rs", "r", encoding="utf-8") as f: + content = f.read() + +git_worker_old = """async fn git_sync_worker(state: Arc) { + let repo_path = std::env::current_dir().unwrap_or_else(|_| ".".into()); + let mut last_commit_id = String::new(); + + loop { + tokio::time::sleep(tokio::time::Duration::from_secs(30)).await; + + if let Ok(repo) = git2::Repository::discover(&repo_path) { + if let Ok(head) = repo.head() { + if let Ok(commit) = head.peel_to_commit() { + let current_id = commit.id().to_string(); + if current_id != last_commit_id && !last_commit_id.is_empty() { + let msg = commit.message().unwrap_or("").to_string(); + let branch = head.shorthand().unwrap_or("unknown").to_string(); + + state.ledger.modify(|changes| { + changes.push(crate::models::CodeChange { + git_commit: Some(current_id.clone()), + git_branch: Some(branch), + description: format!("Auto-synced commit: {}", msg.trim()), + timestamp: std::time::SystemTime::now().duration_since(std::time::UNIX_EPOCH).unwrap().as_secs(), + file_path: "".to_string(), + }); + }); + eprintln!("Git Sync: Logged new commit {}", current_id); + + state.tasks.modify(|tasks| { + for task in tasks.iter_mut() { + if task.status != "completed" && msg.to_lowercase().contains(&task.title.to_lowercase()) { + task.status = "completed".to_string(); + eprintln!("Git Sync: Auto-completed task '{}'", task.title); + } + } + }); + } + last_commit_id = current_id; + } + } + } + } +}""" + +git_worker_new = """async fn git_sync_worker(state: Arc) { + let repo_path = std::env::current_dir().unwrap_or_else(|_| ".".into()); + let mut last_commit_id = String::new(); + + loop { + tokio::time::sleep(tokio::time::Duration::from_secs(30)).await; + + let repo_path_clone = repo_path.clone(); + let commit_data = tokio::task::spawn_blocking(move || { + if let Ok(repo) = git2::Repository::discover(&repo_path_clone) { + if let Ok(head) = repo.head() { + if let Ok(commit) = head.peel_to_commit() { + let current_id = commit.id().to_string(); + let msg = commit.message().unwrap_or("").to_string(); + let branch = head.shorthand().unwrap_or("unknown").to_string(); + return Some((current_id, msg, branch)); + } + } + } + None + }) + .await + .unwrap_or(None); + + if let Some((current_id, msg, branch)) = commit_data { + if current_id != last_commit_id && !last_commit_id.is_empty() { + state.ledger.modify(|changes| { + changes.push(crate::models::CodeChange { + git_commit: Some(current_id.clone()), + git_branch: Some(branch), + description: format!("Auto-synced commit: {}", msg.trim()), + timestamp: std::time::SystemTime::now().duration_since(std::time::UNIX_EPOCH).unwrap().as_secs(), + file_path: "".to_string(), + }); + }); + eprintln!("Git Sync: Logged new commit {}", current_id); + + state.tasks.modify(|tasks| { + for task in tasks.iter_mut() { + if task.status != "completed" && msg.to_lowercase().contains(&task.title.to_lowercase()) { + task.status = "completed".to_string(); + eprintln!("Git Sync: Auto-completed task '{}'", task.title); + } + } + }); + } + last_commit_id = current_id; + } + } +}""" + +content = content.replace(git_worker_old, git_worker_new) + +with open("server/src/main.rs", "w", encoding="utf-8") as f: + f.write(content) diff --git a/fix_git2.py b/fix_git2.py new file mode 100644 index 0000000..d4cfc46 --- /dev/null +++ b/fix_git2.py @@ -0,0 +1,12 @@ +import sys + +cargo_path = r"C:\Users\reazul.ashraf\workspace\rust\mcp-memory\server\Cargo.toml" +with open(cargo_path, "r", encoding="utf-8") as f: + content = f.read() + +# Replace git2 to use vendored openssl to avoid cross-compilation linking issues +content = content.replace('git2 = "0.19.0"', 'git2 = { version = "0.19.0", features = ["vendored-openssl"] }') + +with open(cargo_path, "w", encoding="utf-8") as f: + f.write(content) +print("Added vendored-openssl feature to git2") diff --git a/fix_handlers.py b/fix_handlers.py new file mode 100644 index 0000000..67a9465 --- /dev/null +++ b/fix_handlers.py @@ -0,0 +1,15 @@ +import sys + +handlers_path = r"C:\Users\reazul.ashraf\workspace\rust\mcp-memory\server\src\handlers.rs" +with open(handlers_path, "r", encoding="utf-8") as f: + content = f.read() + +content = content.replace( + 'return Some(crate::mcp::success(req_id, true, &format!("Invalid args: {}", e))),', + 'return Some(crate::mcp::success(id.clone(), serde_json::json!({"isError": true, "content": [{"type": "text", "text": format!("Invalid args: {}", e)}] })));' +) +content = content.replace('let graph = state.graph.read();', 'let graph = self.state.graph.read();') + +with open(handlers_path, "w", encoding="utf-8") as f: + f.write(content) +print("Fixed handlers.rs") diff --git a/fix_index.py b/fix_index.py new file mode 100644 index 0000000..4b9fd44 --- /dev/null +++ b/fix_index.py @@ -0,0 +1,76 @@ +import re + +with open("server/src/state.rs", "r", encoding="utf-8") as f: + content = f.read() + +rebuild_index_old = """ pub fn rebuild_index(&self) { + if let Ok(new_idx) = MemoryIndex::new(&self.base_dir) { + let session = self.session_graph.read().unwrap(); + let mut full = { + let cache = self.master_cache.read().unwrap(); + cache.0.clone() + }; + Self::merge_graphs(&mut full, &session); + for e in full.entities.values() { + let _ = new_idx.index_entity(e); + } + for t in self.tasks.read() { + let _ = new_idx.index_task(&t); + } + for s in self.snippets.read() { + let _ = new_idx.index_snippet(&s); + } + for a in self.adrs.read() { + let _ = new_idx.index_adr(&a); + } + if let Ok(mut w) = self.search_index.write() { + *w = new_idx; + } + } + }""" + +rebuild_index_new = """ pub fn rebuild_index(&self) { + if let Ok(new_idx) = MemoryIndex::new(&self.base_dir) { + let session = self.session_graph.read().unwrap(); + let cache = self.master_cache.read().unwrap(); + + // Index entities that are only in master, or merge if they are in both + for (name, e) in &cache.0.entities { + if let Some(session_e) = session.entities.get(name) { + let mut merged_e = e.clone(); + for obs in &session_e.observations { + if !merged_e.observations.contains(obs) { + merged_e.observations.push(obs.clone()); + } + } + let _ = new_idx.index_entity(&merged_e); + } else { + let _ = new_idx.index_entity(e); + } + } + // Index entities that are only in session + for (name, session_e) in &session.entities { + if !cache.0.entities.contains_key(name) { + let _ = new_idx.index_entity(session_e); + } + } + + for t in self.tasks.read() { + let _ = new_idx.index_task(&t); + } + for s in self.snippets.read() { + let _ = new_idx.index_snippet(&s); + } + for a in self.adrs.read() { + let _ = new_idx.index_adr(&a); + } + if let Ok(mut w) = self.search_index.write() { + *w = new_idx; + } + } + }""" + +content = content.replace(rebuild_index_old, rebuild_index_new) + +with open("server/src/state.rs", "w", encoding="utf-8") as f: + f.write(content) diff --git a/fix_main.py b/fix_main.py index 19f15c1..9240a64 100644 --- a/fix_main.py +++ b/fix_main.py @@ -1,43 +1,37 @@ -import re import sys -main_rs_path = r"C:\Users\reazul.ashraf\workspace\rust\mcp-memory\server\src\main.rs" - -with open(main_rs_path, "r", encoding="utf-8") as f: +main_path = r"C:\Users\reazul.ashraf\workspace\rust\mcp-memory\server\src\main.rs" +with open(main_path, "r", encoding="utf-8") as f: content = f.read() -# Add the route -route_code = """ - .route("/api/tasks/:id/complete", post({ - let state_clone = app_state.handler.state.clone(); - move |axum::extract::Path(id): axum::extract::Path| async move { - if let Ok(mut tasks) = state_clone.tasks.write() { - let mut found = false; - for t in tasks.iter_mut() { - if t.id == id { - t.status = "completed".to_string(); - found = true; - break; - } - } - if found { - let _ = state_clone.flush_store(); - return axum::Json(serde_json::json!({"status": "success"})); - } - } - axum::Json(serde_json::json!({"status": "not_found"})) - } - })) -""" +# Remove WS code +ws_start = content.find('async fn ws_handler') +if ws_start != -1: + content = content[:ws_start] -# Find where to insert it (after /api/tasks) -if ".route(\"/api/tasks/:id/complete\"" not in content: - content = content.replace( - '.route("/api/tasks", get({', - route_code + '\n .route("/api/tasks", get({' - ) +# Remove /ws route +route_idx = content.find('.route("/ws", get(ws_handler))') +if route_idx != -1: + content = content.replace('.route("/ws", get(ws_handler))\n ', '') -with open(main_rs_path, "w", encoding="utf-8") as f: +# Remove WS import +import_idx = content.find('use axum::extract::ws::{WebSocketUpgrade, WebSocket, Message};\n') +if import_idx != -1: + content = content.replace('use axum::extract::ws::{WebSocketUpgrade, WebSocket, Message};\n', '') + +# Fix GC task state +gc_start = content.find(' // Background Garbage Collection for old tasks') +if gc_start != -1: + content = content.replace('let state_gc = state_clone.clone();', 'let state_gc = Arc::clone(&state);') + +# Fix Git task state and CodeChange fields +git_start = content.find(' // Git Native Sync Background Task') +if git_start != -1: + content = content.replace('let state_git = state_clone.clone();', 'let state_git = Arc::clone(&state);') + content = content.replace('commit_hash: current_id.clone(),', 'git_commit: Some(current_id.clone()),') + content = content.replace('branch: branch,', 'git_branch: Some(branch),') + content = content.replace('files_changed: vec![],', 'file_path: "".to_string(),') + +with open(main_path, "w", encoding="utf-8") as f: f.write(content) - -print("Injected /api/tasks/:id/complete route successfully.") +print("Cleaned up WS and fixed main.rs") diff --git a/fix_merge.py b/fix_merge.py new file mode 100644 index 0000000..a486241 --- /dev/null +++ b/fix_merge.py @@ -0,0 +1,52 @@ +import re + +with open("server/src/state.rs", "r", encoding="utf-8") as f: + content = f.read() + +unique_items_block = re.search(r' pub fn unique_items<.*?>\(.*?\) -> Vec \{\n(?:.*?\n){1,15}? \}\n', content, re.DOTALL) +if unique_items_block: + content = content.replace(unique_items_block.group(0), "") + +merge_graphs_old = """ pub fn merge_graphs(dest: &mut KnowledgeGraph, src: &KnowledgeGraph) { + for (name, src_ent) in &src.entities { + let dest_ent = dest + .entities + .entry(name.clone()) + .or_insert_with(|| src_ent.clone()); + if dest_ent.name == src_ent.name { + dest_ent.observations.extend(src_ent.observations.clone()); + dest_ent.observations = Self::unique_items(dest_ent.observations.clone()); + } + } + dest.relations.extend(src.relations.clone()); + dest.relations = Self::unique_items(dest.relations.clone()); + }""" + +merge_graphs_new = """ pub fn merge_graphs(dest: &mut KnowledgeGraph, src: &KnowledgeGraph) { + for (name, src_ent) in &src.entities { + let dest_ent = dest + .entities + .entry(name.clone()) + .or_insert_with(|| crate::models::Entity { + name: src_ent.name.clone(), + entity_type: src_ent.entity_type.clone(), + observations: Vec::new(), + }); + + for obs in &src_ent.observations { + if !dest_ent.observations.contains(obs) { + dest_ent.observations.push(obs.clone()); + } + } + } + for rel in &src.relations { + if !dest.relations.contains(rel) { + dest.relations.push(rel.clone()); + } + } + }""" + +content = content.replace(merge_graphs_old, merge_graphs_new) + +with open("server/src/state.rs", "w", encoding="utf-8") as f: + f.write(content) diff --git a/fix_move.py b/fix_move.py new file mode 100644 index 0000000..edf55ef --- /dev/null +++ b/fix_move.py @@ -0,0 +1,14 @@ +import sys + +main_path = r"C:\Users\reazul.ashraf\workspace\rust\mcp-memory\server\src\main.rs" +with open(main_path, "r", encoding="utf-8") as f: + content = f.read() + +content = content.replace( + 'handler: Arc::new(MemoryHandler { state }),', + 'handler: Arc::new(MemoryHandler { state: Arc::clone(&state) }),' +) + +with open(main_path, "w", encoding="utf-8") as f: + f.write(content) +print("Fixed main.rs move error") diff --git a/fix_reconcile.py b/fix_reconcile.py new file mode 100644 index 0000000..cb7cd43 --- /dev/null +++ b/fix_reconcile.py @@ -0,0 +1,56 @@ +import re + +with open("server/src/main.rs", "r", encoding="utf-8") as f: + content = f.read() + +reconcile_worker_old = """async fn reconcile_worker(state: Arc) { + loop { + sleep(Duration::from_secs(5)).await; + let pattern = format!("{}/delta_*.json", state.base_dir.display()); + let has_local = { + let session = state.session_graph.read().unwrap(); + !session.entities.is_empty() || !session.relations.is_empty() + }; + let has_files = glob::glob(&pattern).map(|p| p.count() > 0).unwrap_or(false); + if has_local || has_files { + state.apply_sync_write(|_master| {}).await; + let state_clone = state.clone(); + let _ = tokio::task::spawn_blocking(move || { + state_clone.rebuild_index(); + }).await; + } + + } +}""" + +reconcile_worker_new = """async fn reconcile_worker(state: Arc) { + loop { + sleep(Duration::from_secs(5)).await; + + let has_local = { + let session = state.session_graph.read().unwrap(); + !session.entities.is_empty() || !session.relations.is_empty() + }; + + let base_dir = state.base_dir.clone(); + let has_files = tokio::task::spawn_blocking(move || { + let pattern = format!("{}/delta_*.json", base_dir.display()); + glob::glob(&pattern).map(|p| p.count() > 0).unwrap_or(false) + }) + .await + .unwrap_or(false); + + if has_local || has_files { + state.apply_sync_write(|_master| {}).await; + let state_clone = state.clone(); + let _ = tokio::task::spawn_blocking(move || { + state_clone.rebuild_index(); + }).await; + } + } +}""" + +content = content.replace(reconcile_worker_old, reconcile_worker_new) + +with open("server/src/main.rs", "w", encoding="utf-8") as f: + f.write(content) diff --git a/fix_route2.py b/fix_route2.py new file mode 100644 index 0000000..2576fbd --- /dev/null +++ b/fix_route2.py @@ -0,0 +1,12 @@ +import re + +main_rs_path = r"C:\Users\reazul.ashraf\workspace\rust\mcp-memory\server\src\main.rs" + +with open(main_rs_path, "r", encoding="utf-8") as f: + content = f.read() + +content = content.replace('.route("/api/tasks/:id/complete"', '.route("/api/tasks/{id}/complete"') + +with open(main_rs_path, "w", encoding="utf-8") as f: + f.write(content) +print("Replaced route syntax.") diff --git a/fix_search_commit.py b/fix_search_commit.py new file mode 100644 index 0000000..5f13f63 --- /dev/null +++ b/fix_search_commit.py @@ -0,0 +1,19 @@ +import re + +with open("server/src/search.rs", "r", encoding="utf-8") as f: + content = f.read() + +# Add pub fn commit +commit_code = """ + pub fn commit(&self) -> tantivy::Result<()> { + let mut writer = self.writer.lock().unwrap(); + writer.commit()?; + Ok(()) + } +""" + +content = content.replace(" pub fn search(", commit_code + "\n pub fn search(") +content = content.replace("let mut writer = self.writer.lock().unwrap();", "let writer = self.writer.lock().unwrap();") + +with open("server/src/search.rs", "w", encoding="utf-8") as f: + f.write(content) diff --git a/fix_state.py b/fix_state.py new file mode 100644 index 0000000..6db6a6e --- /dev/null +++ b/fix_state.py @@ -0,0 +1,21 @@ +import re + +with open("server/src/state.rs", "r", encoding="utf-8") as f: + content = f.read() + +# Add git_branch and namespace to Entity initialization +content = content.replace("observations: Vec::new(),", "observations: Vec::new(),\n namespace: src_ent.namespace.clone(),\n git_branch: src_ent.git_branch.clone(),") + +# Add unique_items back +unique_items_code = """ + pub fn unique_items(input: Vec) -> Vec { + let mut keys = std::collections::HashSet::new(); + input.into_iter().filter(|entry| keys.insert(entry.clone())).collect() + } +""" + +content = content.replace("impl MemoryState {\n", "impl MemoryState {\n" + unique_items_code) + +with open("server/src/state.rs", "w", encoding="utf-8") as f: + f.write(content) + diff --git a/fix_tantivy.py b/fix_tantivy.py new file mode 100644 index 0000000..3815f09 --- /dev/null +++ b/fix_tantivy.py @@ -0,0 +1,36 @@ +import re + +with open("server/src/search.rs", "r", encoding="utf-8") as f: + content = f.read() + +# Remove writer.commit() from index_* +content = content.replace(" writer.commit()?;\n", "") + +# Add pub fn commit +commit_code = """ + pub fn commit(&self) -> tantivy::Result<()> { + let mut writer = self.writer.lock().unwrap(); + writer.commit()?; + Ok(()) + } +""" + +content = content.replace(" pub fn search(&self, query_str: &str, limit: usize) -> tantivy::Result> {", commit_code + "\n pub fn search(&self, query_str: &str, limit: usize) -> tantivy::Result> {") + +with open("server/src/search.rs", "w", encoding="utf-8") as f: + f.write(content) + +with open("server/src/state.rs", "r", encoding="utf-8") as f: + state_content = f.read() + +# Call new_idx.commit() in rebuild_index before locking the main rwlock +state_content = state_content.replace(""" if let Ok(mut w) = self.search_index.write() { + *w = new_idx; + }""", """ let _ = new_idx.commit(); + if let Ok(mut w) = self.search_index.write() { + *w = new_idx; + }""") + +with open("server/src/state.rs", "w", encoding="utf-8") as f: + f.write(state_content) + diff --git a/fix_tools.py b/fix_tools.py new file mode 100644 index 0000000..98197a5 --- /dev/null +++ b/fix_tools.py @@ -0,0 +1,11 @@ +import sys + +tools_path = r"C:\Users\reazul.ashraf\workspace\rust\mcp-memory\server\src\tools.rs" +with open(tools_path, "r", encoding="utf-8") as f: + content = f.read() + +content = content.replace("use schemars::JsonSchema;\nuse serde::{Deserialize, Serialize};\n\n#[derive(Debug, Clone, Serialize, Deserialize, JsonSchema)]\npub struct QueryGraphPathTool", "#[derive(Debug, Clone, Serialize, Deserialize, JsonSchema)]\npub struct QueryGraphPathTool") + +with open(tools_path, "w", encoding="utf-8") as f: + f.write(content) +print("Fixed tools.rs") diff --git a/fix_wal.py b/fix_wal.py new file mode 100644 index 0000000..48e85bd --- /dev/null +++ b/fix_wal.py @@ -0,0 +1,55 @@ +import re + +with open("server/src/state.rs", "r", encoding="utf-8") as f: + content = f.read() + +recover_wal_code = """ + pub fn recover_wal(&self) { + let wal_path = self.base_dir.join("wal.jsonl"); + if let Ok(content) = std::fs::read_to_string(&wal_path) { + let mut session = self.session_graph.write().unwrap(); + for line in content.lines() { + if let Ok(d) = serde_json::from_str::(line) { + Self::merge_graphs(&mut session, &d); + } + } + } + } +""" +content = content.replace(" pub fn merge_graphs(dest: &mut KnowledgeGraph, src: &KnowledgeGraph) {", recover_wal_code + " pub fn merge_graphs(dest: &mut KnowledgeGraph, src: &KnowledgeGraph) {") + +get_full_graph_old = """ pub fn get_full_graph(&self) -> KnowledgeGraph { + let mut master = self.read_master_cached(); + let wal_path = self.base_dir.join("wal.jsonl"); + if let Ok(content) = std::fs::read_to_string(&wal_path) { + for line in content.lines() { + if let Ok(d) = serde_json::from_str::(line) { + Self::merge_graphs(&mut master, &d); + } + } + } + let session_graph = self.session_graph.read().unwrap(); + Self::merge_graphs(&mut master, &session_graph); + master + }""" + +get_full_graph_new = """ pub fn get_full_graph(&self) -> KnowledgeGraph { + let mut master = self.read_master_cached(); + let session_graph = self.session_graph.read().unwrap(); + Self::merge_graphs(&mut master, &session_graph); + master + }""" + +content = content.replace(get_full_graph_old, get_full_graph_new) + +with open("server/src/state.rs", "w", encoding="utf-8") as f: + f.write(content) + +with open("server/src/main.rs", "r", encoding="utf-8") as f: + main_content = f.read() + +main_content = main_content.replace(" state.rebuild_index();\n\n run_server(state)", " state.recover_wal();\n state.rebuild_index();\n\n run_server(state)") + +with open("server/src/main.rs", "w", encoding="utf-8") as f: + f.write(main_content) + diff --git a/fix_writer.py b/fix_writer.py new file mode 100644 index 0000000..a2715ef --- /dev/null +++ b/fix_writer.py @@ -0,0 +1,9 @@ +import re + +with open("server/src/search.rs", "r", encoding="utf-8") as f: + content = f.read() + +content = content.replace("let writer = self.writer.lock().unwrap();\n writer.commit()?;", "let mut writer = self.writer.lock().unwrap();\n writer.commit()?;") + +with open("server/src/search.rs", "w", encoding="utf-8") as f: + f.write(content) diff --git a/inject_gc.py b/inject_gc.py new file mode 100644 index 0000000..13c4ae7 --- /dev/null +++ b/inject_gc.py @@ -0,0 +1,45 @@ +import sys + +main_path = r"C:\Users\reazul.ashraf\workspace\rust\mcp-memory\server\src\main.rs" +with open(main_path, "r", encoding="utf-8") as f: + content = f.read() + +idx = content.find('eprintln!("MCP Memory Server running on http://127.0.0.1:3000/sse");') +if idx == -1: + print("Could not find start point") + sys.exit(1) + +gc_code = """ + // Background Garbage Collection for old tasks + let state_gc = state_clone.clone(); + tokio::spawn(async move { + loop { + // Run every 24 hours + tokio::time::sleep(tokio::time::Duration::from_secs(24 * 3600)).await; + + let now = std::time::SystemTime::now().duration_since(std::time::UNIX_EPOCH).unwrap().as_secs(); + let fourteen_days = 14 * 24 * 3600; + let cutoff = now.saturating_sub(fourteen_days); + + state_gc.tasks.modify(|tasks| { + let initial_len = tasks.len(); + tasks.retain(|task| { + if task.status.to_lowercase() == "completed" && task.created_at < cutoff { + false // remove + } else { + true // keep + } + }); + if tasks.len() < initial_len { + eprintln!("GC: Removed {} old completed tasks", initial_len - tasks.len()); + } + }); + } + }); +""" + +content = content[:idx] + gc_code + content[idx:] +with open(main_path, "w", encoding="utf-8") as f: + f.write(content) + +print("Injected GC task!") diff --git a/inject_git.py b/inject_git.py new file mode 100644 index 0000000..49dbaaf --- /dev/null +++ b/inject_git.py @@ -0,0 +1,62 @@ +import sys + +main_path = r"C:\Users\reazul.ashraf\workspace\rust\mcp-memory\server\src\main.rs" +with open(main_path, "r", encoding="utf-8") as f: + content = f.read() + +idx = content.find('eprintln!("MCP Memory Server running on http://127.0.0.1:3000/sse");') +if idx == -1: + print("Could not find start point") + sys.exit(1) + +git_code = """ + // Git Native Sync Background Task + let state_git = state_clone.clone(); + tokio::spawn(async move { + let repo_path = std::env::current_dir().unwrap_or_else(|_| ".".into()); + let mut last_commit_id = String::new(); + + loop { + tokio::time::sleep(tokio::time::Duration::from_secs(30)).await; + + if let Ok(repo) = git2::Repository::discover(&repo_path) { + if let Ok(head) = repo.head() { + if let Ok(commit) = head.peel_to_commit() { + let current_id = commit.id().to_string(); + if current_id != last_commit_id && !last_commit_id.is_empty() { + let msg = commit.message().unwrap_or("").to_string(); + let branch = head.shorthand().unwrap_or("unknown").to_string(); + + state_git.code_changes.modify(|changes| { + changes.push(crate::models::CodeChange { + commit_hash: current_id.clone(), + branch: branch, + description: format!("Auto-synced commit: {}", msg.trim()), + timestamp: std::time::SystemTime::now().duration_since(std::time::UNIX_EPOCH).unwrap().as_secs(), + files_changed: vec![], + }); + }); + eprintln!("Git Sync: Logged new commit {}", current_id); + + state_git.tasks.modify(|tasks| { + for task in tasks.iter_mut() { + if task.status != "completed" && msg.to_lowercase().contains(&task.title.to_lowercase()) { + task.status = "completed".to_string(); + eprintln!("Git Sync: Auto-completed task '{}'", task.title); + } + } + }); + } + last_commit_id = current_id; + } + } + } + } + }); +""" + +content = content[:idx] + git_code + content[idx:] +with open(main_path, "w", encoding="utf-8") as f: + f.write(content) + +print("Injected Git Sync task!") diff --git a/inject_graph_path.py b/inject_graph_path.py new file mode 100644 index 0000000..82f3192 --- /dev/null +++ b/inject_graph_path.py @@ -0,0 +1,92 @@ +import re +import sys + +handlers_path = r"C:\Users\reazul.ashraf\workspace\rust\mcp-memory\server\src\handlers.rs" + +with open(handlers_path, "r", encoding="utf-8") as f: + content = f.read() + +# Add to list_tools +list_tools_idx = content.find('crate::mcp::tool_def::') +if list_tools_idx == -1: + print("Could not find list_tools") + sys.exit(1) + +tool_def = ' crate::mcp::tool_def::("query_graph_path", "Traverse the knowledge graph to find a path between two entities."),' +content = content[:list_tools_idx] + tool_def + "\n" + content[list_tools_idx:] + +# Add to match name +match_idx = content.find('"create_entities" => {') +if match_idx == -1: + print("Could not find match name") + sys.exit(1) + +handler_code = """ "query_graph_path" => { + let req: crate::tools::QueryGraphPathTool = match parse_args(args.clone()) { + Ok(r) => r, + Err(e) => return Some(crate::mcp::success(req_id, true, &format!("Invalid args: {}", e))), + }; + let graph = state.graph.read(); + let max_depth = req.max_depth.unwrap_or(5); + let mut queue = std::collections::VecDeque::new(); + let mut visited = std::collections::HashSet::new(); + let mut parents: std::collections::HashMap = std::collections::HashMap::new(); + + queue.push_back(req.start_node.clone()); + visited.insert(req.start_node.clone()); + + let mut found = false; + let mut current_depth = 0; + let mut nodes_at_current_depth = 1; + let mut nodes_at_next_depth = 0; + + while let Some(current) = queue.pop_front() { + if current == req.end_node { + found = true; + break; + } + nodes_at_current_depth -= 1; + if current_depth < max_depth { + for rel in &graph.relations { + if rel.from == current && !visited.contains(&rel.to) { + visited.insert(rel.to.clone()); + parents.insert(rel.to.clone(), (current.clone(), rel.relation_type.clone())); + queue.push_back(rel.to.clone()); + nodes_at_next_depth += 1; + } else if rel.to == current && !visited.contains(&rel.from) { + visited.insert(rel.from.clone()); + parents.insert(rel.from.clone(), (current.clone(), format!("inverse({})", rel.relation_type))); + queue.push_back(rel.from.clone()); + nodes_at_next_depth += 1; + } + } + } + if nodes_at_current_depth == 0 { + current_depth += 1; + nodes_at_current_depth = nodes_at_next_depth; + nodes_at_next_depth = 0; + } + } + + if found { + let mut path = Vec::new(); + let mut curr = req.end_node.clone(); + while curr != req.start_node { + let (parent, rel) = parents.get(&curr).unwrap(); + path.push(format!("({}) --[{}]--> ({})", parent, rel, curr)); + curr = parent.clone(); + } + path.reverse(); + Ok(format!("Path found:\\n{}", path.join("\\n"))) + } else { + Ok(format!("No path found between {} and {} within depth {}", req.start_node, req.end_node, max_depth)) + } + } +""" + +content = content[:match_idx] + handler_code + content[match_idx:] + +with open(handlers_path, "w", encoding="utf-8") as f: + f.write(content) + +print("Injected query_graph_path tool!") diff --git a/inject_ws.py b/inject_ws.py new file mode 100644 index 0000000..0d4c2b4 --- /dev/null +++ b/inject_ws.py @@ -0,0 +1,70 @@ +import sys + +main_path = r"C:\Users\reazul.ashraf\workspace\rust\mcp-memory\server\src\main.rs" +with open(main_path, "r", encoding="utf-8") as f: + content = f.read() + +# 1. Add WebSocket imports +import_idx = content.find('use axum::{') +if import_idx == -1: + print("Could not find axum imports") + sys.exit(1) + +import_code = 'use axum::extract::ws::{WebSocketUpgrade, WebSocket, Message};\n' +content = content[:import_idx] + import_code + content[import_idx:] + +# 2. Add /ws route +route_idx = content.find('.route("/sse", get(sse_handler))') +if route_idx == -1: + print("Could not find route definition") + sys.exit(1) + +route_code = '.route("/ws", get(ws_handler))\n ' +content = content[:route_idx] + route_code + content[route_idx:] + +# 3. Add ws_handler implementation +end_idx = len(content) + +ws_code = """ +async fn ws_handler( + ws: WebSocketUpgrade, + State(state): State>, +) -> impl IntoResponse { + ws.on_upgrade(move |socket| handle_socket(socket, state)) +} + +async fn handle_socket(mut socket: WebSocket, state: Arc) { + let mut rx = state.tx.subscribe(); + + // Create a background task to receive messages from the system and send to the WebSocket + let mut send_task = tokio::spawn(async move { + while let Ok(msg) = rx.recv().await { + if socket.send(Message::Text(msg.into())).await.is_err() { + break; + } + } + }); + + // Create a task to receive messages from the WebSocket (e.g. task completion from UI) + // In a real app we'd decode JSON-RPC here + let mut recv_task = tokio::spawn(async move { + // Just keeping it alive and listening + // We can add logic to process JSON-RPC from the frontend here later + loop { + tokio::time::sleep(tokio::time::Duration::from_secs(3600)).await; + } + }); + + tokio::select! { + _ = (&mut send_task) => recv_task.abort(), + _ = (&mut recv_task) => send_task.abort(), + } +} +""" + +content += ws_code + +with open(main_path, "w", encoding="utf-8") as f: + f.write(content) + +print("Injected WS handler!") diff --git a/server/Cargo.toml b/server/Cargo.toml index 3041566..e584a65 100644 --- a/server/Cargo.toml +++ b/server/Cargo.toml @@ -5,13 +5,15 @@ edition = "2024" [dependencies] async-trait = "0.1.92" -axum = "0.8" +axum = { version = "0.8", features = ["ws"] } bincode = "1.3.3" clap = { version = "4.6.6", features = ["derive"] } dashmap = "6.2.1" dirs = "6.0.0" futures-util = "0.3.34" glob = "0.3.4" +git2 = { version = "0.19.0", default-features = false } +notify = "6.1.1" redb = "4.2.0" reqwest = { version = "0.12", default-features = false, features = ["stream", "rustls-tls"] } schemars = "1.2.2" @@ -24,6 +26,7 @@ tokio-util = { version = "0.7.19", features = ["io"] } tracing = "0.1.44" tracing-subscriber = "0.3.23" uuid = { version = "1.26.0", features = ["v4"] } +tokio-tungstenite = "0.21.0" [build-dependencies] chrono = "0.4.45" diff --git a/server/src/dashboard.html b/server/src/dashboard.html index fb4bdf4..701c174 100644 --- a/server/src/dashboard.html +++ b/server/src/dashboard.html @@ -710,32 +710,43 @@ } } - // --- SSE Activity Feed --- - function setupSSE() { - const es = new EventSource('/sse'); + // --- WebSocket Activity Feed --- + function setupWS() { + const protocol = location.protocol === 'https:' ? 'wss:' : 'ws:'; + const ws = new WebSocket(`${protocol}//${location.host}/ws?client=ui`); const feed = document.getElementById('activity-feed'); - es.addEventListener('activity', function(event) { - if (event.data.startsWith('/messages?')) return; - - const div = document.createElement('div'); - div.className = 'feed-entry'; - const time = new Date().toLocaleTimeString(); - div.innerHTML = `[${time}] ${event.data}`; - feed.appendChild(div); - - // Auto-scroll logic - const isScrolledToBottom = feed.scrollHeight - feed.clientHeight <= feed.scrollTop + 20; - if (isScrolledToBottom) { - feed.scrollTop = feed.scrollHeight; + ws.onmessage = function(event) { + try { + const data = JSON.parse(event.data); + if (data.type === 'activity') { + const div = document.createElement('div'); + div.className = 'feed-entry'; + const time = new Date().toLocaleTimeString(); + div.innerHTML = `[${time}] ${data.data}`; + feed.appendChild(div); + + // Auto-scroll logic + const isScrolledToBottom = feed.scrollHeight - feed.clientHeight <= feed.scrollTop + 20; + if (isScrolledToBottom) { + feed.scrollTop = feed.scrollHeight; + } + } + } catch (e) { + // Ignore non-JSON or other messages for now } - }); + }; + + ws.onclose = function() { + console.log("WebSocket closed, attempting to reconnect in 3s..."); + setTimeout(setupWS, 3000); + }; } // --- Start --- loadGraph(); loadTasks(); - setupSSE(); + setupWS(); // Listen to theme toggle changes to update graph font colors const observer = new MutationObserver(() => updateGraphData()); diff --git a/server/src/handlers.rs b/server/src/handlers.rs index ad34cb1..83eb893 100644 --- a/server/src/handlers.rs +++ b/server/src/handlers.rs @@ -39,7 +39,8 @@ impl MemoryHandler { } "tools/list" => { let tools = vec![ - crate::mcp::tool_def::("create_entities", "Create new entities in the knowledge graph."), + crate::mcp::tool_def::("query_graph_path", "Traverse the knowledge graph to find a path between two entities."), +crate::mcp::tool_def::("create_entities", "Create new entities in the knowledge graph."), crate::mcp::tool_def::("create_relations", "Create new relations between entities in the knowledge graph."), crate::mcp::tool_def::("add_observations", "Add new observations to existing entities in the knowledge graph."), crate::mcp::tool_def::("delete_entities", "Delete entities from the knowledge graph."), @@ -120,7 +121,68 @@ impl MemoryHandler { .unwrap_or(serde_json::Value::Object(Default::default())); let result: Result = match name { - "create_entities" => { + "query_graph_path" => { + let req: crate::tools::QueryGraphPathTool = match parse_args(args.clone()) { + Ok(r) => r, + Err(e) => return Some(crate::mcp::success(id.clone(), serde_json::json!({"isError": true, "content": [{"type": "text", "text": format!("Invalid args: {}", e)}] }))) + }; + let graph = self.state.get_full_graph(); + let max_depth = req.max_depth.unwrap_or(5); + let mut queue = std::collections::VecDeque::new(); + let mut visited = std::collections::HashSet::new(); + let mut parents: std::collections::HashMap = std::collections::HashMap::new(); + + queue.push_back(req.start_node.clone()); + visited.insert(req.start_node.clone()); + + let mut found = false; + let mut current_depth = 0; + let mut nodes_at_current_depth = 1; + let mut nodes_at_next_depth = 0; + + while let Some(current) = queue.pop_front() { + if current == req.end_node { + found = true; + break; + } + nodes_at_current_depth -= 1; + if current_depth < max_depth { + for rel in &graph.relations { + if rel.from == current && !visited.contains(&rel.to) { + visited.insert(rel.to.clone()); + parents.insert(rel.to.clone(), (current.clone(), rel.relation_type.clone())); + queue.push_back(rel.to.clone()); + nodes_at_next_depth += 1; + } else if rel.to == current && !visited.contains(&rel.from) { + visited.insert(rel.from.clone()); + parents.insert(rel.from.clone(), (current.clone(), format!("inverse({})", rel.relation_type))); + queue.push_back(rel.from.clone()); + nodes_at_next_depth += 1; + } + } + } + if nodes_at_current_depth == 0 { + current_depth += 1; + nodes_at_current_depth = nodes_at_next_depth; + nodes_at_next_depth = 0; + } + } + + if found { + let mut path = Vec::new(); + let mut curr = req.end_node.clone(); + while curr != req.start_node { + let (parent, rel) = parents.get(&curr).unwrap(); + path.push(format!("({}) --[{}]--> ({})", parent, rel, curr)); + curr = parent.clone(); + } + path.reverse(); + Ok(format!("Path found:\n{}", path.join("\n"))) + } else { + Ok(format!("No path found between {} and {} within depth {}", req.start_node, req.end_node, max_depth)) + } + } +"create_entities" => { let req: CreateEntitiesTool = match parse_args(args.clone()) { Ok(r) => r, Err(e) => { @@ -131,19 +193,12 @@ impl MemoryHandler { } }; self.state.write_to_local_delta(|g| { - for e_val in req.entities { - if let Some(e) = if e_val.is_string() { - serde_json::from_str(e_val.as_str().unwrap()).ok() - } else { - serde_json::from_value(e_val).ok() - } { - let entity: Entity = e; - if !entity.name.is_empty() { - if let Ok(idx) = self.state.search_index.read() { - let _ = idx.index_entity(&entity); - } - g.entities.insert(entity.name.clone(), entity); + for entity in req.entities { + if !entity.name.is_empty() { + if let Ok(idx) = self.state.search_index.read() { + let _ = idx.index_entity(&entity); } + g.entities.insert(entity.name.clone(), entity); } } }).await; @@ -160,16 +215,9 @@ impl MemoryHandler { } }; self.state.write_to_local_delta(|g| { - for r_val in req.relations { - if let Some(r) = if r_val.is_string() { - serde_json::from_str(r_val.as_str().unwrap()).ok() - } else { - serde_json::from_value(r_val).ok() - } { - let relation: Relation = r; - if !relation.from.is_empty() && !relation.to.is_empty() { - g.relations.push(relation); - } + for relation in req.relations { + if !relation.from.is_empty() && !relation.to.is_empty() { + g.relations.push(relation); } } }).await; @@ -185,20 +233,10 @@ impl MemoryHandler { )); } }; - #[derive(Deserialize)] - struct ObsInput { - #[serde(rename = "entityName")] - entity_name: String, - contents: Vec, - } let full = self.state.get_full_graph(); self.state.write_to_local_delta(|g| { - for o_val in req.observations { - if let Some(o) = if o_val.is_string() { - serde_json::from_str::(o_val.as_str().unwrap()).ok() - } else { - serde_json::from_value(o_val).ok() - } && let Some(full_e) = full.entities.get(&o.entity_name) + for o in req.observations { + if let Some(full_e) = full.entities.get(&o.entity_name) { let mut e = g.entities.get(&o.entity_name).cloned().unwrap_or_else( @@ -248,19 +286,9 @@ impl MemoryHandler { )); } }; - #[derive(Deserialize)] - struct ObsDel { - #[serde(rename = "entityName")] - entity_name: String, - observations: Vec, - } self.state.apply_sync_write(|master| { - for d_val in req.deletions { - if let Some(d) = if d_val.is_string() { - serde_json::from_str::(d_val.as_str().unwrap()).ok() - } else { - serde_json::from_value(d_val).ok() - } && let Some(e) = master.entities.get_mut(&d.entity_name) + for d in req.deletions { + if let Some(e) = master.entities.get_mut(&d.entity_name) { let to_rem: HashSet<_> = d.observations.into_iter().collect(); e.observations.retain(|o| !to_rem.contains(o)); @@ -281,17 +309,11 @@ impl MemoryHandler { }; self.state.apply_sync_write(|master| { let mut to_rem = HashSet::new(); - for r_val in req.relations { - if let Some(r) = if r_val.is_string() { - serde_json::from_str::(r_val.as_str().unwrap()).ok() - } else { - serde_json::from_value(r_val).ok() - } { - to_rem.insert(format!( - "{}|{}|{}|{}", - r.from, r.to, r.relation_type, r.namespace - )); - } + for r in req.relations { + to_rem.insert(format!( + "{}|{}|{}|{}", + r.from, r.to, r.relation_type, r.namespace + )); } master.relations.retain(|r| { !to_rem.contains(&format!( diff --git a/server/src/models.rs b/server/src/models.rs index 04dacff..70e2708 100644 --- a/server/src/models.rs +++ b/server/src/models.rs @@ -1,5 +1,6 @@ use serde::{Deserialize, Serialize}; use std::collections::HashMap; +use schemars::JsonSchema; #[derive(Debug, Clone, Serialize, Deserialize, Default)] pub struct CodeChange { @@ -17,7 +18,7 @@ pub struct StickyNote { pub fn default_namespace() -> String { "global".to_string() } -#[derive(Debug, Clone, Serialize, Deserialize)] +#[derive(Debug, Clone, Serialize, Deserialize, JsonSchema)] pub struct Entity { pub name: String, #[serde(rename = "entityType")] @@ -29,7 +30,7 @@ pub struct Entity { #[serde(default)] pub git_branch: Option, } -#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq, Hash)] +#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq, Hash, JsonSchema)] pub struct Relation { pub from: String, pub to: String, diff --git a/server/src/proxy.rs b/server/src/proxy.rs deleted file mode 100644 index e366780..0000000 --- a/server/src/proxy.rs +++ /dev/null @@ -1,121 +0,0 @@ -use futures_util::StreamExt; -use std::sync::Arc; -use tokio::io::AsyncBufReadExt; -use tokio::sync::RwLock; -use tokio::sync::mpsc; -use tokio_util::io::StreamReader; - -pub fn run_proxy(target_url: &str) -> Result> { - let rt = tokio::runtime::Runtime::new().unwrap(); - rt.block_on(async { - let (msg_tx, mut msg_rx) = mpsc::channel::(100); - let (shutdown_tx, mut shutdown_rx) = mpsc::channel::<()>(1); - - tokio::task::spawn_blocking(move || { - let stdin = std::io::stdin(); - let mut handle = stdin.lock(); - let mut buffer = String::new(); - while let Ok(bytes) = std::io::BufRead::read_line(&mut handle, &mut buffer) { - if bytes == 0 { - break; - } - let _ = msg_tx.blocking_send(buffer.clone()); - buffer.clear(); - } - let _ = shutdown_tx.blocking_send(()); - }); - - let target_url = target_url.to_string(); - let post_url = Arc::new(RwLock::new(String::new())); - let post_url_clone = Arc::clone(&post_url); - - let client = reqwest::Client::builder().build().unwrap(); - - tokio::spawn(async move { - while let Some(msg) = msg_rx.recv().await { - loop { - let url = post_url_clone.read().await.clone(); - if !url.is_empty() { - let res = client - .post(&url) - .header("Accept", "application/json, text/event-stream") - .header("Content-Type", "application/json") - .body(msg.clone()) - .send() - .await; - - if res.is_ok() { - break; - } - } - tokio::time::sleep(tokio::time::Duration::from_millis(100)).await; - } - } - }); - - if shutdown_rx.try_recv().is_ok() { - return Ok(false); - } - - let sse_url = format!("{}/sse", target_url); - let client = reqwest::Client::builder().build().unwrap(); - - match client - .get(&sse_url) - .header("Accept", "text/event-stream") - .send() - .await - { - Ok(resp) => { - if resp.status() == reqwest::StatusCode::GONE { - return Ok(true); - } - - let stream = resp.bytes_stream().map(|res| { - res.map_err(std::io::Error::other) - }); - let mut reader = tokio::io::BufReader::new(StreamReader::new(stream)); - let mut line = String::new(); - let mut is_message = false; - let mut is_endpoint = false; - - loop { - tokio::select! { - _ = shutdown_rx.recv() => { - return Ok(false); - } - res = reader.read_line(&mut line) => { - match res { - Ok(bytes) => { - if bytes == 0 { break; } - let trimmed = line.trim(); - if trimmed.starts_with("event: message") { - is_message = true; - is_endpoint = false; - } else if trimmed.starts_with("event: endpoint") { - is_endpoint = true; - is_message = false; - } else if let Some(stripped) = trimmed.strip_prefix("data: ") { - if is_message { - println!("{}", stripped); - is_message = false; - } else if is_endpoint { - let mut p = post_url.write().await; - *p = format!("{}{}", target_url, stripped); - is_endpoint = false; - } - } - line.clear(); - } - Err(_) => break, - } - } - } - } - *post_url.write().await = String::new(); - Ok(true) - } - Err(_) => Ok(true), - } - }) -} diff --git a/server/src/tools.rs b/server/src/tools.rs index 4282373..4cdd57c 100644 --- a/server/src/tools.rs +++ b/server/src/tools.rs @@ -1,276 +1,516 @@ use schemars::JsonSchema; use serde::{Deserialize, Serialize}; +/// Create new entities in the knowledge graph. #[derive(Debug, Deserialize, Serialize, JsonSchema)] pub struct CreateEntitiesTool { - pub entities: Vec, + /// Array of entities to create. + pub entities: Vec, } + +/// Create new relations between entities in the knowledge graph. #[derive(Debug, Deserialize, Serialize, JsonSchema)] pub struct CreateRelationsTool { - pub relations: Vec, + /// Array of relations to create. + pub relations: Vec, } + +#[derive(Debug, Deserialize, Serialize, JsonSchema)] +pub struct ObservationInput { + #[serde(rename = "entityName")] + pub entity_name: String, + pub contents: Vec, +} + +/// Add new observations to existing entities in the knowledge graph. #[derive(Debug, Deserialize, Serialize, JsonSchema)] pub struct AddObservationsTool { - pub observations: Vec, + /// Array of observations to add. + pub observations: Vec, } + +/// Delete entities from the knowledge graph. #[derive(Debug, Deserialize, Serialize, JsonSchema)] pub struct DeleteEntitiesTool { + /// Array of entity names to delete. #[serde(rename = "entityNames")] pub entity_names: Vec, } + #[derive(Debug, Deserialize, Serialize, JsonSchema)] -pub struct DeleteObservationsTool { - pub deletions: Vec, -} -#[derive(Debug, Deserialize, Serialize, JsonSchema)] -pub struct DeleteRelationsTool { - pub relations: Vec, -} -#[derive(Debug, Deserialize, Serialize, JsonSchema)] -pub struct ReadGraphTool { - pub namespace: Option, -} -#[derive(Debug, Deserialize, Serialize, JsonSchema)] -pub struct SearchNodesTool { - pub query: String, - pub namespace: Option, -} -#[derive(Debug, Deserialize, Serialize, JsonSchema)] -pub struct OpenNodesTool { - pub names: Vec, -} -#[derive(Debug, Deserialize, Serialize, JsonSchema)] -pub struct LogCodeChangeTool { - #[serde(rename = "filePath")] - pub file_path: String, - pub description: String, - pub git_commit: Option, - pub git_branch: Option, -} -#[derive(Debug, Deserialize, Serialize, JsonSchema)] -pub struct QueryRecentChangesTool {} -#[derive(Debug, Deserialize, Serialize, JsonSchema)] -pub struct VisualizeGraphTool { - pub query: Option, - pub namespace: Option, -} -#[derive(Debug, Deserialize, Serialize, JsonSchema)] -pub struct AddStickyNoteTool { - pub content: String, -} -#[derive(Debug, Deserialize, Serialize, JsonSchema)] -pub struct ReadStickyNotesTool {} -#[derive(Debug, Deserialize, Serialize, JsonSchema)] -pub struct CondenseEntityTool { +pub struct DeleteObservationInput { #[serde(rename = "entityName")] pub entity_name: String, + pub observations: Vec, +} + +/// Delete observations from existing entities. +#[derive(Debug, Deserialize, Serialize, JsonSchema)] +pub struct DeleteObservationsTool { + /// Array of observation deletions. + pub deletions: Vec, +} + +/// Delete relations between entities. +#[derive(Debug, Deserialize, Serialize, JsonSchema)] +pub struct DeleteRelationsTool { + /// Array of relations to delete. + pub relations: Vec, +} + +/// Read the entire knowledge graph. +#[derive(Debug, Deserialize, Serialize, JsonSchema)] +pub struct ReadGraphTool { + /// Optional namespace to restrict the read to. + pub namespace: Option, +} + +/// Search for entities in the knowledge graph by name or type. +#[derive(Debug, Deserialize, Serialize, JsonSchema)] +pub struct SearchNodesTool { + /// The search query. + pub query: String, + /// Optional namespace to restrict the search to. + pub namespace: Option, +} + +/// Open and retrieve full details of specific nodes in the knowledge graph. +#[derive(Debug, Deserialize, Serialize, JsonSchema)] +pub struct OpenNodesTool { + /// Array of entity names to open. + pub names: Vec, +} + +/// Log a significant code change or refactor in the memory system. +#[derive(Debug, Deserialize, Serialize, JsonSchema)] +pub struct LogCodeChangeTool { + /// The path of the file that was changed. + #[serde(rename = "filePath")] + pub file_path: String, + /// A description of the change. + pub description: String, + /// The associated git commit hash, if any. + pub git_commit: Option, + /// The associated git branch, if any. + pub git_branch: Option, +} + +/// Query recently logged code changes. +#[derive(Debug, Deserialize, Serialize, JsonSchema)] +pub struct QueryRecentChangesTool {} + +/// Generate a visual representation of the knowledge graph. +#[derive(Debug, Deserialize, Serialize, JsonSchema)] +pub struct VisualizeGraphTool { + /// Optional search query to filter the graph before visualization. + pub query: Option, + /// Optional namespace to restrict the visualization to. + pub namespace: Option, +} + +/// Add a sticky note for unstructured thoughts or reminders. +#[derive(Debug, Deserialize, Serialize, JsonSchema)] +pub struct AddStickyNoteTool { + /// The content of the sticky note. + pub content: String, +} + +/// Read all active sticky notes. +#[derive(Debug, Deserialize, Serialize, JsonSchema)] +pub struct ReadStickyNotesTool {} + +/// Condense or summarize an entity's observations to reduce size. +#[derive(Debug, Deserialize, Serialize, JsonSchema)] +pub struct CondenseEntityTool { + /// The name of the entity to condense. + #[serde(rename = "entityName")] + pub entity_name: String, + /// The condensed observations that will replace the existing ones. pub summarized_observations: Vec, } + +/// Add a new task to the task tracker. #[derive(Debug, Deserialize, Serialize, JsonSchema)] pub struct AddTaskTool { + /// The title of the task. pub title: String, + /// A detailed description of the task. pub description: String, + /// The associated git branch, if any. pub git_branch: Option, } + +/// Update the status of an existing task. #[derive(Debug, Deserialize, Serialize, JsonSchema)] pub struct UpdateTaskStatusTool { + /// The ID of the task to update. pub id: String, + /// The new status of the task (e.g., 'pending' or 'completed'). + #[schemars(description = "Must be 'pending' or 'completed'")] pub status: String, } + +/// List all currently active tasks. #[derive(Debug, Deserialize, Serialize, JsonSchema)] pub struct ListActiveTasksTool { + /// Optional git branch to filter tasks by. pub git_branch: Option, } + +/// Store a reusable code snippet. #[derive(Debug, Deserialize, Serialize, JsonSchema)] pub struct StoreSnippetTool { + /// The name of the snippet. pub name: String, + /// The programming language of the snippet. pub language: String, + /// The code snippet content. pub code: String, + /// A description of what the snippet does. pub description: String, } + +/// Search through stored code snippets. #[derive(Debug, Deserialize, Serialize, JsonSchema)] pub struct SearchSnippetsTool { + /// The search query. pub query: String, } + +/// Delete a stored code snippet. #[derive(Debug, Deserialize, Serialize, JsonSchema)] pub struct DeleteSnippetTool { + /// The name of the snippet to delete. pub name: String, } + +/// Log an architectural decision record (ADR). #[derive(Debug, Deserialize, Serialize, JsonSchema)] pub struct LogDecisionTool { + /// The title of the decision. pub title: String, + /// The context or problem requiring a decision. pub context: String, + /// The decision made. pub decision: String, + /// The consequence of the decision. pub consequence: String, } + +/// Query architectural decision records. #[derive(Debug, Deserialize, Serialize, JsonSchema)] pub struct QueryDecisionsTool { + /// Optional search query. pub query: Option, } + +/// Merge two entities in the knowledge graph into one. #[derive(Debug, Deserialize, Serialize, JsonSchema)] pub struct MergeEntitiesTool { + /// The name of the entity to merge from (will be deleted). #[serde(rename = "sourceEntity")] pub source_entity: String, + /// The name of the entity to merge into. #[serde(rename = "targetEntity")] pub target_entity: String, } + +/// Find orphaned entities (entities without any relations) in the graph. #[derive(Debug, Deserialize, Serialize, JsonSchema)] pub struct FindOrphansTool {} + +/// Record a user preference or behavior to adapt future interactions. #[derive(Debug, Deserialize, Serialize, JsonSchema)] pub struct LearnPreferenceTool { + /// The key for the preference. pub key: String, + /// The value of the preference. pub value: String, } + +/// Read all learned user preferences. #[derive(Debug, Deserialize, Serialize, JsonSchema)] pub struct ReadPreferencesTool {} + +/// Log a complex error and its fix for future reference. #[derive(Debug, Deserialize, Serialize, JsonSchema)] pub struct LogErrorFixTool { + /// The error signature or stack trace. pub signature: String, + /// The solution applied to fix the error. pub solution: String, + /// The associated git commit hash, if any. pub git_commit: Option, + /// The associated git branch, if any. pub git_branch: Option, } + +/// Search through previously logged error fixes. #[derive(Debug, Deserialize, Serialize, JsonSchema)] pub struct SearchErrorFixesTool { + /// The search query. pub query: String, } + +/// Pin a file to keep it explicitly in the context workspace. #[derive(Debug, Deserialize, Serialize, JsonSchema)] pub struct PinFileTool { + /// The namespace to pin the file in. pub namespace: String, + /// The absolute path of the file to pin. pub file_path: String, + /// The associated git branch, if any. pub git_branch: Option, } + +/// Unpin a file from the context workspace. #[derive(Debug, Deserialize, Serialize, JsonSchema)] pub struct UnpinFileTool { + /// The namespace the file is pinned in. pub namespace: String, + /// The absolute path of the file to unpin. pub file_path: String, } + +/// List all currently pinned files. #[derive(Debug, Deserialize, Serialize, JsonSchema)] pub struct ListPinnedFilesTool { + /// Optional namespace to filter by. pub namespace: Option, + /// Optional git branch to filter by. pub git_branch: Option, } + +/// Add a summary of the current session. #[derive(Debug, Deserialize, Serialize, JsonSchema)] pub struct AddSessionSummaryTool { + /// The summary content. pub summary: String, + /// The namespace to add the summary to. pub namespace: String, } + +/// Get a timeline of major project events. #[derive(Debug, Deserialize, Serialize, JsonSchema)] pub struct GetProjectTimelineTool { + /// Optional namespace to restrict the timeline to. pub namespace: Option, } + +/// Leave a memo for the next session or agent. #[derive(Debug, Deserialize, Serialize, JsonSchema)] pub struct LeaveHandoffMemoTool { + /// The content of the memo. pub content: String, + /// The namespace to leave the memo in. pub namespace: String, } + +/// Read pending handoff memos. #[derive(Debug, Deserialize, Serialize, JsonSchema)] pub struct ReadHandoffMemosTool { + /// Optional namespace to restrict the read to. pub namespace: Option, } + +/// Clear handoff memos after reading. #[derive(Debug, Deserialize, Serialize, JsonSchema)] pub struct ClearHandoffMemosTool { + /// Array of memo IDs to clear. pub ids: Vec, } + +/// Update the environment fingerprint (e.g., OS, tool versions). #[derive(Debug, Deserialize, Serialize, JsonSchema)] pub struct UpdateEnvFingerprintTool { + /// The namespace to update the fingerprint for. pub namespace: String, + /// A map of tool names to their versions. pub tool_versions: std::collections::HashMap, } + +/// Read the environment fingerprint. #[derive(Debug, Deserialize, Serialize, JsonSchema)] pub struct ReadEnvFingerprintTool { + /// The namespace to read the fingerprint for. pub namespace: String, } + +/// Log a required environment variable or configuration. #[derive(Debug, Deserialize, Serialize, JsonSchema)] pub struct LogEnvRequirementTool { + /// The namespace. pub namespace: String, + /// The environment variable key (e.g., DATABASE_URL). pub key: String, + /// A description of what the variable is used for. pub description: String, + /// Whether the variable contains a secret. pub is_secret: bool, } + +/// Add a new milestone. #[derive(Debug, Deserialize, Serialize, JsonSchema)] pub struct AddMilestoneTool { + /// The title of the milestone. pub title: String, + /// The namespace for the milestone. pub namespace: String, } + +/// Update the status of a milestone. #[derive(Debug, Deserialize, Serialize, JsonSchema)] pub struct UpdateMilestoneTool { + /// The ID of the milestone to update. pub id: String, + /// The new status (e.g., 'active', 'completed'). pub status: String, } + +/// List all milestones. #[derive(Debug, Deserialize, Serialize, JsonSchema)] pub struct ListMilestonesTool { + /// Optional namespace to restrict the list to. pub namespace: Option, } + +/// Generate a standup report for a specific time window. #[derive(Debug, Deserialize, Serialize, JsonSchema)] pub struct GenerateStandupReportTool { + /// The namespace to generate the report for. pub namespace: String, + /// The number of hours to look back for activity. pub hours_lookback: u64, } + +/// Register a new infrastructure environment. #[derive(Debug, Deserialize, Serialize, JsonSchema)] pub struct RegisterEnvironmentTool { + /// The namespace. pub namespace: String, + /// The name of the environment (e.g., 'staging', 'prod'). pub name: String, + /// The URL or connection string for the environment. pub url: String, + /// A description of the environment. pub description: String, + /// Whether a VPN is required to access the environment. pub requires_vpn: bool, } + +/// Get details about a registered environment. #[derive(Debug, Deserialize, Serialize, JsonSchema)] pub struct GetEnvironmentDetailsTool { + /// The namespace to retrieve details for. pub namespace: String, } + +/// Add an item to the PR checklist. #[derive(Debug, Deserialize, Serialize, JsonSchema)] pub struct AddPrChecklistItemTool { + /// The namespace. pub namespace: String, + /// The description of the checklist item. pub description: String, } + +/// Get the PR checklist. #[derive(Debug, Deserialize, Serialize, JsonSchema)] pub struct GetPrChecklistTool { + /// The namespace. pub namespace: String, } + +/// Clear the PR checklist. #[derive(Debug, Deserialize, Serialize, JsonSchema)] pub struct ClearPrChecklistTool { + /// The namespace. pub namespace: String, } + +/// Log a technical debt record. #[derive(Debug, Deserialize, Serialize, JsonSchema)] pub struct LogTechDebtTool { + /// The namespace. pub namespace: String, + /// A description of the technical debt. pub description: String, + /// The ideal solution to resolve the debt. pub ideal_solution: String, + /// The associated git commit hash, if any. pub git_commit: Option, + /// The associated git branch, if any. pub git_branch: Option, } + +/// Resolve a technical debt record. #[derive(Debug, Deserialize, Serialize, JsonSchema)] pub struct ResolveTechDebtTool { + /// The ID of the technical debt record to resolve. pub id: String, } + +/// List technical debt records. #[derive(Debug, Deserialize, Serialize, JsonSchema)] pub struct ListTechDebtTool { + /// The namespace. pub namespace: String, + /// Whether to include resolved technical debt in the results. pub include_resolved: bool, } + +/// Save the current context workspace. #[derive(Debug, Deserialize, Serialize, JsonSchema)] pub struct SaveContextWorkspaceTool { + /// The namespace for the workspace. pub namespace: String, + /// The name to save the workspace as. pub name: String, + /// Array of pinned file paths. pub pinned_files: Vec, + /// Array of active task IDs. pub active_task_ids: Vec, } + +/// Load a saved context workspace. #[derive(Debug, Deserialize, Serialize, JsonSchema)] pub struct LoadContextWorkspaceTool { + /// The namespace of the workspace. pub namespace: String, + /// The name of the workspace to load. pub name: String, } + +/// List all saved context workspaces. #[derive(Debug, Deserialize, Serialize, JsonSchema)] pub struct ListContextWorkspacesTool { + /// The namespace to list workspaces for. pub namespace: String, } + +/// Search across all memory stores (Graph, Tasks, Snippets, ADRs, etc.). #[derive(Debug, Deserialize, Serialize, JsonSchema)] pub struct OmniSearchTool { + /// The search query. pub query: String, + /// Optional namespace to restrict the search to. pub namespace: Option, } + +/// Get a health digest of the project. #[derive(Debug, Deserialize, Serialize, JsonSchema)] pub struct GetProjectHealthTool { + /// The namespace to get health for. pub namespace: String, } + +/// Traverse the knowledge graph to find a path between two entities. +#[derive(Debug, Clone, Serialize, Deserialize, JsonSchema)] +pub struct QueryGraphPathTool { + /// The starting entity name. + pub start_node: String, + /// The ending entity name. + pub end_node: String, + /// Optional maximum depth to search. + pub max_depth: Option, +} diff --git a/server_main_old.rs b/server_main_old.rs new file mode 100644 index 0000000..adc3bfc --- /dev/null +++ b/server_main_old.rs @@ -0,0 +1,823 @@ +mod handlers; +mod mcp; +mod models; +mod search; +mod state; +mod store; +mod tools; + +use crate::handlers::MemoryHandler; +use crate::models::*; +use crate::state::MemoryState; +use crate::store::Store; + +use std::fs; +use std::path::PathBuf; +use std::sync::{Arc, RwLock}; +use std::time::{Duration, SystemTime, UNIX_EPOCH}; +use redb::ReadableTable; +use tokio::time::sleep; + +use clap::{Parser, Subcommand}; +use std::collections::HashMap; + +#[derive(Parser)] +#[command(author, version, about = "Antigravity MCP Memory Server", long_about = None)] +struct Cli { + #[command(subcommand)] + command: Option, + /// Target URL for the proxy to connect to (e.g., http://127.0.0.1:3000) + #[arg(long)] + target: Option, + /// Run the server as a background daemon process (Windows only) + #[arg(long)] + daemon: bool, + /// Send a shutdown request to the currently running server + #[arg(long)] + exit: bool, + /// Send a shutdown request to the existing server and wait for it to exit + #[arg(long)] + restart: bool, +} + +#[derive(Subcommand)] +enum Commands { + /// Manage authorization gates and verification for actions + Gate { + #[command(subcommand)] + subcmd: GateCommands, + }, +} + +#[derive(Subcommand)] +enum GateCommands { + Set { + #[arg(long)] + action: String, + #[arg(long)] + target: String, + #[arg(long)] + namespace: Option, + #[arg(short = 'p', long = "param")] + params: Vec, + #[arg(long, conflicts_with = "block")] + authorize: bool, + #[arg(long, conflicts_with = "authorize")] + block: bool, + #[arg(long)] + reason: Option, + }, + Verify { + #[arg(long)] + action: String, + #[arg(long)] + target: String, + #[arg(long)] + namespace: Option, + #[arg(short = 'p', long = "param")] + params: Vec, + #[arg(long)] + consume: bool, + }, +} + +async fn reconcile_worker(state: Arc) { + loop { + sleep(Duration::from_secs(5)).await; + let pattern = format!("{}/delta_*.json", state.base_dir.display()); + let has_local = { + let session = state.session_graph.read().unwrap(); + !session.entities.is_empty() || !session.relations.is_empty() + }; + let has_files = glob::glob(&pattern).map(|p| p.count() > 0).unwrap_or(false); + if has_local || has_files { + state.apply_sync_write(|_master| {}).await; + let state_clone = state.clone(); + let _ = tokio::task::spawn_blocking(move || { + state_clone.rebuild_index(); + }).await; + } + + let now = SystemTime::now() + .duration_since(UNIX_EPOCH) + .unwrap() + .as_secs(); + state.ledger.modify(|ledger| { + let seven_days = now.saturating_sub(7 * 24 * 60 * 60); + ledger.retain(|c| c.timestamp >= seven_days); + if ledger.len() > 1000 { + let excess = ledger.len() - 1000; + ledger.drain(0..excess); + } + }); + state.sticky.modify(|notes| { + notes.retain(|note| note.timestamp >= now.saturating_sub(24 * 60 * 60)); + }); + } +} + +use axum::{ + Json, Router, + extract::{Query, State, ws::{WebSocketUpgrade, WebSocket, Message}}, + response::IntoResponse, + routing::{get, post}, +}; +use futures_util::{SinkExt, StreamExt}; +use std::sync::atomic::{AtomicUsize, Ordering}; +use tokio::sync::mpsc; + +struct AppState { + handler: Arc, + clients: RwLock>>, + next_id: AtomicUsize, +} + +#[derive(serde::Deserialize)] +struct GateVerifyReq { + action: String, + target: String, + namespace: Option, + #[serde(default)] + params: HashMap, + #[serde(default)] + consume: bool, +} + +#[derive(serde::Deserialize)] +struct GateSetReq { + action: String, + target: String, + namespace: Option, + #[serde(default)] + params: HashMap, + authorize: Option, + block: Option, + reason: Option, +} + +async fn gate_verify_handler( + State(app_state): State>, + Query(q): Query, +) -> axum::response::Response { + let mut found = None; + let mut to_remove = None; + app_state.handler.state.gates.modify(|gates| { + if let Some(idx) = gates.iter().position(|g| { + g.action == q.action + && g.target == q.target + && g.namespace == q.namespace + && g.params == q.params + }) { + found = Some(gates[idx].clone()); + if q.consume { + to_remove = Some(idx); + } + } + if let Some(idx) = to_remove { + gates.remove(idx); + } + }); + + match found { + Some(record) => { + if record.status == "authorized" { + (axum::http::StatusCode::OK, "Authorized").into_response() + } else { + let msg = if let Some(r) = record.reason { + format!("Action blocked. Reason: {}", r) + } else { + "Action blocked.".to_string() + }; + (axum::http::StatusCode::FORBIDDEN, msg).into_response() + } + } + None => { + (axum::http::StatusCode::NOT_FOUND, "Action not yet authorized (no gate record found).").into_response() + } + } +} + +async fn gate_set_handler( + State(app_state): State>, + Json(body): Json, +) -> axum::response::Response { + let status = if body.block.unwrap_or(false) { + "blocked".to_string() + } else if body.authorize.unwrap_or(false) { + "authorized".to_string() + } else { + "pending".to_string() + }; + + let record = GateRecord { + id: uuid::Uuid::new_v4().to_string(), + action: body.action.clone(), + target: body.target.clone(), + namespace: body.namespace.clone(), + params: body.params.clone(), + status, + reason: body.reason.clone(), + timestamp: SystemTime::now() + .duration_since(UNIX_EPOCH) + .unwrap() + .as_secs(), + }; + app_state.handler.state.gates.modify(|gates| { + gates.retain(|g| !(g.action == record.action && g.target == record.target)); + gates.push(record); + }); + (axum::http::StatusCode::OK, "Gate state updated.").into_response() +} + +fn run_server(state: Arc) -> Result<(), Box> { + let rt = tokio::runtime::Runtime::new().unwrap(); + rt.block_on(async { + tokio::spawn(reconcile_worker(Arc::clone(&state))); + let app_state = Arc::new(AppState { + handler: Arc::new(MemoryHandler { state: Arc::clone(&state) }), + clients: RwLock::new(HashMap::new()), + next_id: AtomicUsize::new(1), + }); + + let app = Router::new() + .route("/api/version", get(|| async move { + axum::Json(serde_json::json!({ + "version": env!("BUILD_DATE"), + "git_hash": option_env!("GIT_HASH").unwrap_or("unknown") + })) + })) + .route("/ws", get(ws_handler)) + .route("/health", get(health_handler)) + .route("/gate/verify", get(gate_verify_handler)) + .route("/gate/set", post(gate_set_handler)) + .route( + "/shutdown", + post(|| async move { + std::thread::spawn(|| { + std::thread::sleep(std::time::Duration::from_millis(100)); + std::process::exit(0); + }); + "Shutting down..." + }), + ) + .route( + "/", + get(|| async move { axum::response::Html(include_str!("dashboard.html")) }), + ) + .route("/api/graph", get({ + let state_clone = app_state.handler.state.clone(); + move || async move { + let graph = state_clone.get_full_graph(); + axum::Json(graph) + } + })) + + .route("/api/tasks/{id}/complete", post({ + let state_clone = app_state.handler.state.clone(); + move |axum::extract::Path(id): axum::extract::Path| async move { + state_clone.tasks.modify(|tasks| { + for t in tasks.iter_mut() { + if t.id == id { + t.status = "completed".to_string(); + break; + } + } + }); + axum::Json(serde_json::json!({"status": "success"})) + } + })) + + .route("/api/tasks", get({ + let state_clone = app_state.handler.state.clone(); + move || async move { + let tasks = state_clone.tasks.read(); + axum::Json(tasks.clone()) + } + })) + .route("/api/search", get({ + let state_clone = app_state.handler.state.clone(); + move |axum::extract::Query(params): axum::extract::Query>| async move { + if let Some(q) = params.get("q") { + if let Ok(idx) = state_clone.search_index.read() { + if let Ok(results) = idx.search(q, None) { + let mut formatted_results = Vec::new(); + for (type_name, content) in results { + formatted_results.push(serde_json::json!({ + "type_name": type_name, + "content": content, + "score": 1.0 + })); + } + return axum::Json(serde_json::json!({ "results": formatted_results })); + } + } + } + axum::Json(serde_json::json!({ "results": [] })) + } + })) + .route("/api/stats", + get({ + let state_clone = app_state.handler.state.clone(); + move || async move { + let (entities, relations) = { + let graph = state_clone.get_full_graph(); + (graph.entities.len(), graph.relations.len()) + }; + let tasks = state_clone.tasks.read().len(); + let snippets = state_clone.snippets.read().len(); + let tech_debts = state_clone.tech_debts.read().len(); + let adrs = state_clone.adrs.read().len(); + + let ledger = state_clone.ledger.read().len(); + let sticky = state_clone.sticky.read().len(); + let error_fixes = state_clone.error_fixes.read().len(); + let pinned_files = state_clone.pinned_files.read().len(); + let session_summaries = state_clone.session_summaries.read().len(); + let handoff_memos = state_clone.handoff_memos.read().len(); + let env_fingerprints = state_clone.env_fingerprints.read().len(); + let env_requirements = state_clone.env_requirements.read().len(); + let milestones = state_clone.milestones.read().len(); + let environments = state_clone.environments.read().len(); + let pr_checklists = state_clone.pr_checklists.read().len(); + let gates = state_clone.gates.read().len(); + let context_workspaces = state_clone.context_workspaces.read().len(); + + axum::Json(serde_json::json!({ + "entities": entities, + "relations": relations, + "tasks": tasks, + "snippets": snippets, + "tech_debts": tech_debts, + "adrs": adrs, + "ledger": ledger, + "sticky": sticky, + "error_fixes": error_fixes, + "pinned_files": pinned_files, + "session_summaries": session_summaries, + "handoff_memos": handoff_memos, + "env_fingerprints": env_fingerprints, + "env_requirements": env_requirements, + "milestones": milestones, + "environments": environments, + "pr_checklists": pr_checklists, + "gates": gates, + "context_workspaces": context_workspaces + })) + } + }), + ) + .with_state(app_state); + + let mut retries = 0; + let listener = loop { + match tokio::net::TcpListener::bind("127.0.0.1:3000").await { + Ok(l) => break l, + Err(e) => { + // Check if it's already running and healthy + if let Ok(mut stream) = std::net::TcpStream::connect("127.0.0.1:3000") { + use std::io::{Read, Write}; + let _ = stream.write_all( + b"GET /health HTTP/1.1\r\nHost: 127.0.0.1\r\nConnection: close\r\n\r\n", + ); + let mut response = String::new(); + let _ = stream.read_to_string(&mut response); + if response.contains("200 OK") { + // Already healthy! Just exit cleanly instead of panicking/retrying loop. + std::process::exit(0); + } + } + + retries += 1; + if retries > 15 { + let log_path = dirs::home_dir() + .unwrap_or_default() + .join(".gemini/mcp_memory/daemon_fatal.log"); + let _ = std::fs::write( + &log_path, + format!( + "FATAL: Could not bind to 127.0.0.1:3000 after 15 seconds: {}\n", + e + ), + ); + std::process::exit(1); + } + let log_path = dirs::home_dir() + .unwrap_or_default() + .join(".gemini/mcp_memory/daemon_error.log"); + if let Ok(mut file) = std::fs::OpenOptions::new() + .create(true) + .append(true) + .open(&log_path) + { + use std::io::Write; + let _ = writeln!( + file, + "Failed to bind to 127.0.0.1:3000 (attempt {}): {}. Retrying in 1s...", + retries, e + ); + } + tokio::time::sleep(std::time::Duration::from_secs(1)).await; + } + } + }; + + // Background Garbage Collection for old tasks + let state_gc = Arc::clone(&state); + tokio::spawn(async move { + loop { + // Run every 24 hours + tokio::time::sleep(tokio::time::Duration::from_secs(24 * 3600)).await; + + let now = std::time::SystemTime::now().duration_since(std::time::UNIX_EPOCH).unwrap().as_secs(); + let fourteen_days = 14 * 24 * 3600; + let cutoff = now.saturating_sub(fourteen_days); + + state_gc.tasks.modify(|tasks| { + let initial_len = tasks.len(); + tasks.retain(|task| { + if task.status.to_lowercase() == "completed" && task.created_at < cutoff { + false // remove + } else { + true // keep + } + }); + if tasks.len() < initial_len { + eprintln!("GC: Removed {} old completed tasks", initial_len - tasks.len()); + } + }); + } + }); + + // Git Native Sync Background Task + let state_git = Arc::clone(&state); + tokio::spawn(async move { + let repo_path = std::env::current_dir().unwrap_or_else(|_| ".".into()); + let mut last_commit_id = String::new(); + + loop { + tokio::time::sleep(tokio::time::Duration::from_secs(30)).await; + + if let Ok(repo) = git2::Repository::discover(&repo_path) { + if let Ok(head) = repo.head() { + if let Ok(commit) = head.peel_to_commit() { + let current_id = commit.id().to_string(); + if current_id != last_commit_id && !last_commit_id.is_empty() { + let msg = commit.message().unwrap_or("").to_string(); + let branch = head.shorthand().unwrap_or("unknown").to_string(); + + state_git.ledger.modify(|changes| { + changes.push(crate::models::CodeChange { + git_commit: Some(current_id.clone()), + git_branch: Some(branch), + description: format!("Auto-synced commit: {}", msg.trim()), + timestamp: std::time::SystemTime::now().duration_since(std::time::UNIX_EPOCH).unwrap().as_secs(), + file_path: "".to_string(), + }); + }); + eprintln!("Git Sync: Logged new commit {}", current_id); + + state_git.tasks.modify(|tasks| { + for task in tasks.iter_mut() { + if task.status != "completed" && msg.to_lowercase().contains(&task.title.to_lowercase()) { + task.status = "completed".to_string(); + eprintln!("Git Sync: Auto-completed task '{}'", task.title); + } + } + }); + } + last_commit_id = current_id; + } + } + } + } + }); +eprintln!("MCP Memory Server running on http://127.0.0.1:3000/sse"); + if let Err(e) = axum::serve(listener, app).await { + let log_path = dirs::home_dir() + .unwrap_or_default() + .join(".gemini/mcp_memory/daemon_error.log"); + let _ = std::fs::write(&log_path, format!("Server crashed: {}\n", e)); + } + Ok(()) + }) +} + +async fn ws_handler( + ws: WebSocketUpgrade, + State(state): State>, + Query(query): Query>, +) -> impl axum::response::IntoResponse { + let client_type = query.get("client").cloned().unwrap_or_else(|| "unknown".to_string()); + ws.on_upgrade(move |socket| handle_socket(socket, state, client_type)) +} + +async fn handle_socket(socket: WebSocket, state: Arc, client_type: String) { + let session_id = format!("{}", state.next_id.fetch_add(1, Ordering::SeqCst)); + let (tx, mut rx) = mpsc::channel::(100); + + state.clients.write().unwrap().insert(session_id.clone(), tx.clone()); + + let (mut sender, mut receiver) = socket.split(); + + let mut send_task = tokio::spawn(async move { + while let Some(msg) = rx.recv().await { + if sender.send(Message::Text(msg.into())).await.is_err() { + break; + } + } + }); + + let handler = Arc::clone(&state.handler); + let state_clone = Arc::clone(&state); + let session_id_clone = session_id.clone(); + + let mut recv_task = tokio::spawn(async move { + while let Some(Ok(Message::Text(text))) = receiver.next().await { + if let Ok(payload) = serde_json::from_str::(&text) { + if client_type == "proxy" { + // Send activity broadcast to UI clients + if let Some(method) = payload.get("method").and_then(|m| m.as_str()) { + if method == "tools/call" { + let name = payload.get("params").and_then(|p| p.get("name")).and_then(|n| n.as_str()).unwrap_or("unknown_tool"); + let activity_msg = format!("Agent executed tool: {}", name); + + let event = serde_json::json!({ + "type": "activity", + "data": activity_msg + }); + + let clients_map = state_clone.clients.read().unwrap().clone(); + for (id, client_tx) in clients_map.iter() { + if id != &session_id_clone { + let _ = client_tx.send(event.to_string()).await; + } + } + } + } + + // Process MCP request + if let Some(response) = handler.handle_request(payload).await { + let res_str = serde_json::to_string(&response).unwrap(); + let tx_opt = state_clone.clients.read().unwrap().get(&session_id_clone).cloned(); + if let Some(client_tx) = tx_opt { + let _ = client_tx.send(res_str).await; + } + } + } + } + } + }); + + tokio::select! { + _ = (&mut send_task) => recv_task.abort(), + _ = (&mut recv_task) => send_task.abort(), + }; + + state.clients.write().unwrap().remove(&session_id); +} + +async fn health_handler() -> &'static str { + "OK" +} + + +fn main() -> Result<(), Box> { + let cli = Cli::parse(); + + if cli.exit { + if let Ok(mut stream) = std::net::TcpStream::connect("127.0.0.1:3000") { + use std::io::Write; + let _ = stream.write_all( + b"POST /shutdown HTTP/1.1\r\nHost: 127.0.0.1\r\nConnection: close\r\n\r\n", + ); + } + println!("Sent shutdown request to server."); + return Ok(()); + } + + if cli.restart { + if let Ok(mut stream) = std::net::TcpStream::connect("127.0.0.1:3000") { + use std::io::Write; + let _ = stream.write_all( + b"POST /shutdown HTTP/1.1\r\nHost: 127.0.0.1\r\nConnection: close\r\n\r\n", + ); + println!("Sent shutdown request to existing server. Waiting for it to exit..."); + std::thread::sleep(std::time::Duration::from_millis(1500)); + } + return Ok(()); + } + + #[cfg(target_os = "windows")] + { + use std::os::windows::process::CommandExt; + if !cli.daemon { + // Just spawn the daemon and exit. We no longer act as a proxy. + #[allow(clippy::zombie_processes)] + let _ = std::process::Command::new(std::env::current_exe().unwrap()) + .arg("--daemon") + .stdin(std::process::Stdio::null()) + .stdout(std::process::Stdio::null()) + .stderr(std::process::Stdio::null()) + .creation_flags(0x08000000) // CREATE_NO_WINDOW + .spawn() + .expect("Failed to spawn daemon"); + return Ok(()); + } + } + + #[cfg(not(target_os = "windows"))] + { + // Linux no longer executes server logic natively due to workspace split + return Ok(()); + } + + let base_dir = std::env::var("MCP_MEMORY_STORE_DIR").unwrap_or_else(|_| { + dirs::home_dir() + .map(|mut h| { + h.push(".gemini/mcp_memory"); + h.to_string_lossy().into_owned() + }) + .unwrap_or_else(|| ".gemini/mcp_memory".into()) + }); + let base = PathBuf::from(base_dir); + fs::create_dir_all(&base).expect("Failed to create store dir"); + + let redb_path = base.join("mcp_store.redb"); + let db = Arc::new(redb::Database::create(&redb_path).unwrap()); + + // Ensure table exists and migrate old JSON files + { + let write_txn = db.begin_write().unwrap(); + { + let mut table = write_txn.open_table(crate::store::STORE_TABLE).unwrap(); + + let stores = [ + ("audit_ledger", "audit_ledger.json"), + ("sticky_notes", "sticky_notes.json"), + ("tasks", "tasks.json"), + ("snippets", "snippets.json"), + ("adrs", "adrs.json"), + ("preferences", "preferences.json"), + ("error_fixes", "error_fixes.json"), + ("pinned_files", "pinned_files.json"), + ("session_summaries", "session_summaries.json"), + ("handoff_memos", "handoff_memos.json"), + ("env_fingerprints", "env_fingerprints.json"), + ("env_requirements", "env_requirements.json"), + ("milestones", "milestones.json"), + ("environments", "environments.json"), + ("pr_checklists", "pr_checklists.json"), + ("tech_debts", "tech_debts.json"), + ("gates", "gates.json"), + ("context_workspaces", "context_workspaces.json"), + ]; + + for (key, file_name) in stores.iter() { + if table.get(*key).unwrap().is_none() { + let json_path = base.join(file_name); + if json_path.exists() { + if let Ok(data) = fs::read(&json_path) { + if serde_json::from_slice::(&data).is_ok() { + table.insert(*key, data.as_slice()).unwrap(); + } + } + } + } + } + } + write_txn.commit().unwrap(); + } + + let state = Arc::new(MemoryState { + master_path: base.join("knowledge_graph_master.json"), + session_graph: RwLock::new(KnowledgeGraph::default()), + base_dir: base.clone(), + master_cache: RwLock::new((KnowledgeGraph::default(), SystemTime::UNIX_EPOCH)), + search_index: RwLock::new(crate::search::MemoryIndex::new(&base).unwrap()), + ledger: Store::new("audit_ledger", db.clone()), + sticky: Store::new("sticky_notes", db.clone()), + tasks: Store::new("tasks", db.clone()), + snippets: Store::new("snippets", db.clone()), + adrs: Store::new("adrs", db.clone()), + prefs: Store::new("preferences", db.clone()), + error_fixes: Store::new("error_fixes", db.clone()), + pinned_files: Store::new("pinned_files", db.clone()), + session_summaries: Store::new("session_summaries", db.clone()), + handoff_memos: Store::new("handoff_memos", db.clone()), + env_fingerprints: Store::new("env_fingerprints", db.clone()), + env_requirements: Store::new("env_requirements", db.clone()), + milestones: Store::new("milestones", db.clone()), + environments: Store::new("environments", db.clone()), + pr_checklists: Store::new("pr_checklists", db.clone()), + tech_debts: Store::new("tech_debts", db.clone()), + gates: Store::new("gates", db.clone()), + context_workspaces: Store::new("context_workspaces", db.clone()), + }); + + state.rebuild_index(); + + if let Some(command) = cli.command { + match command { + Commands::Gate { subcmd } => match subcmd { + GateCommands::Set { + action, + target, + namespace, + params, + authorize, + block, + reason, + } => { + let status = if authorize { + "authorized".to_string() + } else if block { + "blocked".to_string() + } else { + "pending".to_string() + }; + let mut param_map = HashMap::new(); + for p in params { + if let Some((k, v)) = p.split_once('=') { + param_map.insert(k.to_string(), v.to_string()); + } + } + let record = GateRecord { + id: uuid::Uuid::new_v4().to_string(), + action: action.clone(), + target: target.clone(), + namespace, + params: param_map, + status, + reason, + timestamp: SystemTime::now() + .duration_since(UNIX_EPOCH) + .unwrap() + .as_secs(), + }; + state.gates.modify(|gates| { + gates.retain(|g| !(g.action == record.action && g.target == record.target)); + gates.push(record); + }); + println!("Gate state updated."); + std::process::exit(0); + } + GateCommands::Verify { + action, + target, + namespace, + params, + consume, + } => { + let mut param_map = HashMap::new(); + for p in params { + if let Some((k, v)) = p.split_once('=') { + param_map.insert(k.to_string(), v.to_string()); + } + } + let mut found = None; + let mut to_remove = None; + state.gates.modify(|gates| { + if let Some(idx) = gates.iter().position(|g| { + g.action == action + && g.target == target + && g.namespace == namespace + && g.params == param_map + }) { + found = Some(gates[idx].clone()); + if consume { + to_remove = Some(idx); + } + } + if let Some(idx) = to_remove { + gates.remove(idx); + } + }); + + match found { + Some(record) => { + if record.status == "authorized" { + std::process::exit(0); + } else { + if let Some(r) = record.reason { + eprintln!("❌ Action blocked. Reason: {}", r); + } else { + eprintln!("❌ Action blocked."); + } + std::process::exit(1); + } + } + None => { + eprintln!("❌ Action not yet authorized (no gate record found)."); + std::process::exit(2); + } + } + } + }, + } + } + + run_server(state) +} + + + diff --git a/strip_openssl.py b/strip_openssl.py new file mode 100644 index 0000000..36003f7 --- /dev/null +++ b/strip_openssl.py @@ -0,0 +1,21 @@ +import sys + +cargo_path = r"C:\Users\reazul.ashraf\workspace\rust\mcp-memory\server\Cargo.toml" +with open(cargo_path, "r", encoding="utf-8") as f: + content = f.read() + +# Swap reqwest rustls-tls in instead of default (which uses native-tls / openssl) +# This entirely removes the openssl dependency during Linux cross-compilation +if 'reqwest = { version = "0.12", default-features = false, features = ["stream", "rustls-tls"] }' not in content: + content = content.replace('reqwest = "0.12"', 'reqwest = { version = "0.12", default-features = false, features = ["stream", "rustls-tls"] }') + +# For git2, we don't need vendored openssl if we're not using https fetch/push for libgit2. +# We just need to read local `.git` files, so we can disable the default features. +if 'git2 = { version = "0.19.0", features = ["vendored-openssl"] }' in content: + content = content.replace('git2 = { version = "0.19.0", features = ["vendored-openssl"] }', 'git2 = { version = "0.19.0", default-features = false }') +elif 'git2 = "0.19.0"' in content: + content = content.replace('git2 = "0.19.0"', 'git2 = { version = "0.19.0", default-features = false }') + +with open(cargo_path, "w", encoding="utf-8") as f: + f.write(content) +print("Removed OpenSSL dependency entirely from server Cargo.toml")