diff --git a/cmd/agentsview/cli.go b/cmd/agentsview/cli.go index f3ba8c59c8..6d81c65347 100644 --- a/cmd/agentsview/cli.go +++ b/cmd/agentsview/cli.go @@ -127,6 +127,7 @@ func newRootCommand() *cobra.Command { root.AddCommand(newClassifierCommand()) root.AddCommand(newSecretsCommand()) root.AddCommand(newSkillsCommand()) + root.AddCommand(newCloudCommand()) root.AddCommand(newDoctorCommand()) root.AddCommand(newVersionCommand()) root.AddCommand(newOpenAPICommand()) diff --git a/cmd/agentsview/cloud.go b/cmd/agentsview/cloud.go new file mode 100644 index 0000000000..a676a2f55f --- /dev/null +++ b/cmd/agentsview/cloud.go @@ -0,0 +1,15 @@ +package main + +import "github.com/spf13/cobra" + +// newCloudCommand is retained as a stable CLI group. Claude.ai connection and +// synchronization are desktop transport features; the Go CLI never reads +// browser credentials or makes provider HTTP requests. +func newCloudCommand() *cobra.Command { + return &cobra.Command{ + Use: "cloud", + Short: "Manage cloud conversation sources", + SilenceUsage: true, + Args: cobra.NoArgs, + } +} diff --git a/desktop/src-tauri/.gitignore b/desktop/src-tauri/.gitignore index 68cb28b4d0..732fa5e553 100644 --- a/desktop/src-tauri/.gitignore +++ b/desktop/src-tauri/.gitignore @@ -10,3 +10,4 @@ # Backup created by prepare-sidecar.sh version patching tauri.conf.json.orig +permissions/autogenerated/ diff --git a/desktop/src-tauri/Cargo.lock b/desktop/src-tauri/Cargo.lock index 458c88bb95..9bf1e60c7b 100644 --- a/desktop/src-tauri/Cargo.lock +++ b/desktop/src-tauri/Cargo.lock @@ -12,6 +12,7 @@ checksum = "320119579fcad9c21884f5c4861d16174d0e06250625266f50fe6898340abefa" name = "agentsview-desktop" version = "0.1.0" dependencies = [ + "serde", "serde_json", "tauri", "tauri-build", @@ -249,6 +250,21 @@ version = "0.22.1" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "72b3254f16251a8381aa12e40e3c4d2f0199f8c6508fbecb9d91f575e0fbb8c6" +[[package]] +name = "bit-set" +version = "0.8.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "08807e080ed7f9d5433fa9b275196cfc35414f66a0c79d864dc51a0d825231a3" +dependencies = [ + "bit-vec", +] + +[[package]] +name = "bit-vec" +version = "0.8.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "5e764a1d40d510daf35e07be9eb06e75770908c27d411ee6c92109c9840eaaf7" + [[package]] name = "bitflags" version = "1.3.2" @@ -524,9 +540,9 @@ checksum = "773648b94d0e5d620f64f280777445740e61fe701025087ec8b57f45c791888b" [[package]] name = "core-graphics" -version = "0.24.0" +version = "0.25.0" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "fa95a34622365fa5bbf40b20b75dba8dfa8c94c734aea8ac9a5ca38af14316f1" +checksum = "064badf302c3194842cf2c5d61f56cc88e54a759313879cdf03abdd27d0c3b97" dependencies = [ "bitflags 2.11.0", "core-foundation", @@ -606,6 +622,19 @@ dependencies = [ "syn 1.0.109", ] +[[package]] +name = "cssparser" +version = "0.36.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "dae61cf9c0abb83bd659dab65b7e4e38d8236824c85f0f804f173567bda257d2" +dependencies = [ + "cssparser-macros", + "dtoa-short", + "itoa", + "phf 0.13.1", + "smallvec", +] + [[package]] name = "cssparser-macros" version = "0.6.1" @@ -618,14 +647,20 @@ dependencies = [ [[package]] name = "ctor" -version = "0.2.9" +version = "0.8.0" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "32a2785755761f3ddc1492979ce1e48d2c00d09311c39e4466429188f3dd6501" +checksum = "352d39c2f7bef1d6ad73db6f5160efcaed66d94ef8c6c573a8410c00bf909a98" dependencies = [ - "quote", - "syn 2.0.117", + "ctor-proc-macro", + "dtor", ] +[[package]] +name = "ctor-proc-macro" +version = "0.0.7" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "52560adf09603e58c9a7ee1fe1dcb95a16927b17c127f0ac02d6e768a0e25bc1" + [[package]] name = "darling" version = "0.21.3" @@ -661,6 +696,17 @@ dependencies = [ "syn 2.0.117", ] +[[package]] +name = "dbus" +version = "0.9.12" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "3ab69f03cc8c4340c9c8e315114e1658e6775a9b16a04357973aa21cec22b32e" +dependencies = [ + "libc", + "libdbus-sys", + "windows-sys 0.61.2", +] + [[package]] name = "deranged" version = "0.3.11" @@ -695,6 +741,27 @@ dependencies = [ "syn 2.0.117", ] +[[package]] +name = "derive_more" +version = "2.1.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "d751e9e49156b02b44f9c1815bcb94b984cdcc4396ecc32521c739452808b134" +dependencies = [ + "derive_more-impl", +] + +[[package]] +name = "derive_more-impl" +version = "2.1.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "799a97264921d8623a957f6c3b9011f3b5492f557bbb7a5a19b7fa6d06ba8dcb" +dependencies = [ + "proc-macro2", + "quote", + "rustc_version", + "syn 2.0.117", +] + [[package]] name = "digest" version = "0.10.7" @@ -726,12 +793,6 @@ dependencies = [ "windows-sys 0.61.2", ] -[[package]] -name = "dispatch" -version = "0.2.0" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "bd0c93bb4b0c6d9b77f4435b0ae98c24d17f1c45b2ff844c6151a07256ca923b" - [[package]] name = "dispatch2" version = "0.3.1" @@ -778,6 +839,21 @@ dependencies = [ "syn 2.0.117", ] +[[package]] +name = "dom_query" +version = "0.27.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "521e380c0c8afb8d9a1e83a1822ee03556fc3e3e7dbc1fd30be14e37f9cb3f89" +dependencies = [ + "bit-set", + "cssparser 0.36.0", + "foldhash 0.2.0", + "html5ever 0.38.0", + "precomputed-hash", + "selectors 0.36.1", + "tendril 0.5.1", +] + [[package]] name = "dpi" version = "0.1.2" @@ -802,6 +878,21 @@ dependencies = [ "dtoa", ] +[[package]] +name = "dtor" +version = "0.3.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "f1057d6c64987086ff8ed0fd3fbf377a6b7d205cc7715868cd401705f715cbe4" +dependencies = [ + "dtor-proc-macro", +] + +[[package]] +name = "dtor-proc-macro" +version = "0.0.6" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "f678cf4a922c215c63e0de95eb1ff08a958a81d47e485cf9da1e27bf6305cfa5" + [[package]] name = "dunce" version = "1.0.5" @@ -982,6 +1073,12 @@ version = "0.1.5" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "d9c4f5dac5e15c24eb999c26181a6ca40b39fe946cbe4c263c7209467bc83af2" +[[package]] +name = "foldhash" +version = "0.2.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "77ce24cb58228fbb8aa041425bb1050850ac19177686ea6e0f41a70416f56fdb" + [[package]] name = "foreign-types" version = "0.5.0" @@ -1437,7 +1534,7 @@ version = "0.15.5" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "9229cfe53dfd69f0609a49f65461bd93001ea1ef889cd5529dd176593f5338a1" dependencies = [ - "foldhash", + "foldhash 0.1.5", ] [[package]] @@ -1478,10 +1575,20 @@ checksum = "3b7410cae13cbc75623c98ac4cbfd1f0bedddf3227afc24f370cf0f50a44a11c" dependencies = [ "log", "mac", - "markup5ever", + "markup5ever 0.14.1", "match_token", ] +[[package]] +name = "html5ever" +version = "0.38.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "1054432bae2f14e0061e33d23402fbaa67a921d319d56adc6bcf887ddad1cbc2" +dependencies = [ + "log", + "markup5ever 0.38.0", +] + [[package]] name = "http" version = "1.4.0" @@ -1909,18 +2016,12 @@ version = "0.8.8-speedreader" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "02cb977175687f33fa4afa0c95c112b987ea1443e5a51c8f8ff27dc618270cc2" dependencies = [ - "cssparser", - "html5ever", + "cssparser 0.29.6", + "html5ever 0.29.1", "indexmap 2.13.0", - "selectors", + "selectors 0.24.0", ] -[[package]] -name = "lazy_static" -version = "1.5.0" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "bbd2bcb4c963f2ddae06a2efc7e9f3591312473c50c6685e1f298068316e66fe" - [[package]] name = "leb128fmt" version = "0.1.0" @@ -1957,6 +2058,15 @@ version = "0.2.182" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "6800badb6cb2082ffd7b6a67e6125bb39f18782f793520caee8cb8846be06112" +[[package]] +name = "libdbus-sys" +version = "0.2.7" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "328c4789d42200f1eeec05bd86c9c13c7f091d2ba9a6ea35acdf51f31bc0f043" +dependencies = [ + "pkg-config", +] + [[package]] name = "libloading" version = "0.7.4" @@ -2021,9 +2131,20 @@ dependencies = [ "log", "phf 0.11.3", "phf_codegen 0.11.3", - "string_cache", - "string_cache_codegen", - "tendril", + "string_cache 0.8.9", + "string_cache_codegen 0.5.4", + "tendril 0.4.3", +] + +[[package]] +name = "markup5ever" +version = "0.38.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "8983d30f2915feeaaab2d6babdd6bc7e9ed1a00b66b5e6d74df19aa9c0e91862" +dependencies = [ + "log", + "tendril 0.5.1", + "web_atoms", ] [[package]] @@ -2103,9 +2224,9 @@ dependencies = [ [[package]] name = "muda" -version = "0.17.1" +version = "0.19.3" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "01c1738382f66ed56b3b9c8119e794a2e23148ac8ea214eda86622d4cb9d415a" +checksum = "1dd04e60bc0b07438a6771710ee1698f98f6ebbc7f89b61264af1563b8aeb878" dependencies = [ "crossbeam-channel", "dpi", @@ -2116,10 +2237,10 @@ dependencies = [ "objc2-core-foundation", "objc2-foundation", "once_cell", - "png 0.17.16", + "png 0.18.1", "serde", "thiserror 2.0.18", - "windows-sys 0.60.2", + "windows-sys 0.61.2", ] [[package]] @@ -2137,12 +2258,6 @@ dependencies = [ "thiserror 1.0.69", ] -[[package]] -name = "ndk-context" -version = "0.1.1" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "27b02d87554356db9e9a873add8782d4ea6e3e58ea071a9adb9a2e8ddb884a8b" - [[package]] name = "ndk-sys" version = "0.6.0+11769913" @@ -2219,17 +2334,9 @@ checksum = "d49e936b501e5c5bf01fda3a9452ff86dc3ea98ad5f283e1455153142d97518c" dependencies = [ "bitflags 2.11.0", "block2", - "libc", "objc2", - "objc2-cloud-kit", - "objc2-core-data", "objc2-core-foundation", - "objc2-core-graphics", - "objc2-core-image", - "objc2-core-text", - "objc2-core-video", "objc2-foundation", - "objc2-quartz-core", ] [[package]] @@ -2249,7 +2356,6 @@ version = "0.3.2" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "0b402a653efbb5e82ce4df10683b6b28027616a2715e90009947d50b8dd298fa" dependencies = [ - "bitflags 2.11.0", "objc2", "objc2-foundation", ] @@ -2289,28 +2395,25 @@ dependencies = [ ] [[package]] -name = "objc2-core-text" +name = "objc2-core-location" version = "0.3.2" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "0cde0dfb48d25d2b4862161a4d5fcc0e3c24367869ad306b0c9ec0073bfed92d" +checksum = "ca347214e24bc973fc025fd0d36ebb179ff30536ed1f80252706db19ee452009" dependencies = [ - "bitflags 2.11.0", "objc2", - "objc2-core-foundation", - "objc2-core-graphics", + "objc2-foundation", ] [[package]] -name = "objc2-core-video" +name = "objc2-core-text" version = "0.3.2" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "d425caf1df73233f29fd8a5c3e5edbc30d2d4307870f802d18f00d83dc5141a6" +checksum = "0cde0dfb48d25d2b4862161a4d5fcc0e3c24367869ad306b0c9ec0073bfed92d" dependencies = [ "bitflags 2.11.0", "objc2", "objc2-core-foundation", "objc2-core-graphics", - "objc2-io-surface", ] [[package]] @@ -2352,16 +2455,6 @@ dependencies = [ "objc2-core-foundation", ] -[[package]] -name = "objc2-javascript-core" -version = "0.3.2" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "2a1e6550c4caed348956ce3370c9ffeca70bb1dbed4fa96112e7c6170e074586" -dependencies = [ - "objc2", - "objc2-core-foundation", -] - [[package]] name = "objc2-osa-kit" version = "0.3.2" @@ -2387,25 +2480,33 @@ dependencies = [ ] [[package]] -name = "objc2-security" +name = "objc2-ui-kit" version = "0.3.2" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "709fe137109bd1e8b5a99390f77a7d8b2961dafc1a1c5db8f2e60329ad6d895a" +checksum = "d87d638e33c06f577498cbcc50491496a3ed4246998a7fbba7ccb98b1e7eab22" dependencies = [ "bitflags 2.11.0", + "block2", "objc2", + "objc2-cloud-kit", + "objc2-core-data", "objc2-core-foundation", + "objc2-core-graphics", + "objc2-core-image", + "objc2-core-location", + "objc2-core-text", + "objc2-foundation", + "objc2-quartz-core", + "objc2-user-notifications", ] [[package]] -name = "objc2-ui-kit" +name = "objc2-user-notifications" version = "0.3.2" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "d87d638e33c06f577498cbcc50491496a3ed4246998a7fbba7ccb98b1e7eab22" +checksum = "9df9128cbbfef73cda168416ccf7f837b62737d748333bfe9ab71c245d76613e" dependencies = [ - "bitflags 2.11.0", "objc2", - "objc2-core-foundation", "objc2-foundation", ] @@ -2421,8 +2522,6 @@ dependencies = [ "objc2-app-kit", "objc2-core-foundation", "objc2-foundation", - "objc2-javascript-core", - "objc2-security", ] [[package]] @@ -2581,10 +2680,20 @@ version = "0.11.3" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "1fd6780a80ae0c52cc120a26a1a42c1ae51b247a253e4e06113d23d2c2edd078" dependencies = [ - "phf_macros 0.11.3", "phf_shared 0.11.3", ] +[[package]] +name = "phf" +version = "0.13.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "c1562dc717473dbaa4c1f85a36410e03c047b2e7df7f45ee938fbef64ae7fadf" +dependencies = [ + "phf_macros 0.13.1", + "phf_shared 0.13.1", + "serde", +] + [[package]] name = "phf_codegen" version = "0.8.0" @@ -2605,6 +2714,16 @@ dependencies = [ "phf_shared 0.11.3", ] +[[package]] +name = "phf_codegen" +version = "0.13.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "49aa7f9d80421bca176ca8dbfebe668cc7a2684708594ec9f3c0db0805d5d6e1" +dependencies = [ + "phf_generator 0.13.1", + "phf_shared 0.13.1", +] + [[package]] name = "phf_generator" version = "0.8.0" @@ -2635,6 +2754,16 @@ dependencies = [ "rand 0.8.5", ] +[[package]] +name = "phf_generator" +version = "0.13.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "135ace3a761e564ec88c03a77317a7c6b80bb7f7135ef2544dbe054243b89737" +dependencies = [ + "fastrand", + "phf_shared 0.13.1", +] + [[package]] name = "phf_macros" version = "0.10.0" @@ -2651,12 +2780,12 @@ dependencies = [ [[package]] name = "phf_macros" -version = "0.11.3" +version = "0.13.1" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "f84ac04429c13a7ff43785d75ad27569f2951ce0ffd30a3321230db2fc727216" +checksum = "812f032b54b1e759ccd5f8b6677695d5268c588701effba24601f6932f8269ef" dependencies = [ - "phf_generator 0.11.3", - "phf_shared 0.11.3", + "phf_generator 0.13.1", + "phf_shared 0.13.1", "proc-macro2", "quote", "syn 2.0.117", @@ -2689,6 +2818,15 @@ dependencies = [ "siphasher 1.0.2", ] +[[package]] +name = "phf_shared" +version = "0.13.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "e57fef6bc5981e38c2ce2d63bfa546861309f875b8a75f092d1d54ae2d64f266" +dependencies = [ + "siphasher 1.0.2", +] + [[package]] name = "pin-project-lite" version = "0.2.17" @@ -3157,6 +3295,12 @@ dependencies = [ "windows-sys 0.52.0", ] +[[package]] +name = "rustc-hash" +version = "2.1.3" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "6b1e7f9a428571be2dc5bc0505c13fb6bf936822b894ec87abf8a08a4e51742d" + [[package]] name = "rustc_version" version = "0.4.1" @@ -3363,14 +3507,33 @@ source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "0c37578180969d00692904465fb7f6b3d50b9a2b952b87c23d0e2e5cb5013416" dependencies = [ "bitflags 1.3.2", - "cssparser", - "derive_more", + "cssparser 0.29.6", + "derive_more 0.99.20", "fxhash", "log", "phf 0.8.0", "phf_codegen 0.8.0", "precomputed-hash", - "servo_arc", + "servo_arc 0.2.0", + "smallvec", +] + +[[package]] +name = "selectors" +version = "0.36.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "c5d9c0c92a92d33f08817311cf3f2c29a3538a8240e94a6a3c622ce652d7e00c" +dependencies = [ + "bitflags 2.11.0", + "cssparser 0.36.0", + "derive_more 2.1.1", + "log", + "new_debug_unreachable", + "phf 0.13.1", + "phf_codegen 0.13.1", + "precomputed-hash", + "rustc-hash", + "servo_arc 0.4.3", "smallvec", ] @@ -3542,6 +3705,15 @@ dependencies = [ "stable_deref_trait", ] +[[package]] +name = "servo_arc" +version = "0.4.3" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "170fb83ab34de17dc69aa7c67482b22218ddb85da56546f9bd6b929e32a05930" +dependencies = [ + "stable_deref_trait", +] + [[package]] name = "sha2" version = "0.10.9" @@ -3708,6 +3880,18 @@ dependencies = [ "serde", ] +[[package]] +name = "string_cache" +version = "0.9.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "a18596f8c785a729f2819c0f6a7eae6ebeebdfffbfe4214ae6b087f690e31901" +dependencies = [ + "new_debug_unreachable", + "parking_lot", + "phf_shared 0.13.1", + "precomputed-hash", +] + [[package]] name = "string_cache_codegen" version = "0.5.4" @@ -3720,6 +3904,18 @@ dependencies = [ "quote", ] +[[package]] +name = "string_cache_codegen" +version = "0.6.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "585635e46db231059f76c5849798146164652513eb9e8ab2685939dd90f29b69" +dependencies = [ + "phf_generator 0.13.1", + "phf_shared 0.13.1", + "proc-macro2", + "quote", +] + [[package]] name = "strsim" version = "0.11.1" @@ -3800,35 +3996,35 @@ dependencies = [ [[package]] name = "tao" -version = "0.34.5" +version = "0.35.3" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "f3a753bdc39c07b192151523a3f77cd0394aa75413802c883a0f6f6a0e5ee2e7" +checksum = "d1c93047acf68669466a34690ac58cca7010bd1b201e1ec86f1fd0a75d3dd4a9" dependencies = [ "bitflags 2.11.0", "block2", "core-foundation", "core-graphics", "crossbeam-channel", - "dispatch", + "dbus", + "dispatch2", "dlopen2", "dpi", "gdkwayland-sys", "gdkx11-sys", "gtk", "jni", - "lazy_static", "libc", "log", "ndk", - "ndk-context", "ndk-sys", "objc2", "objc2-app-kit", "objc2-foundation", + "objc2-ui-kit", "once_cell", "parking_lot", + "percent-encoding", "raw-window-handle", - "scopeguard", "tao-macros", "unicode-segmentation", "url", @@ -3868,9 +4064,9 @@ checksum = "61c41af27dd6d1e27b1b16b489db798443478cef1f06a660c96db617ba5de3b1" [[package]] name = "tauri" -version = "2.10.2" +version = "2.11.1" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "463ae8677aa6d0f063a900b9c41ecd4ac2b7ca82f0b058cc4491540e55b20129" +checksum = "b93bd86d231f0a8138f11a02a584769fe4b703dc36ae133d783228dbc4801405" dependencies = [ "anyhow", "bytes", @@ -3920,9 +4116,9 @@ dependencies = [ [[package]] name = "tauri-build" -version = "2.5.5" +version = "2.6.3" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "ca7bd893329425df750813e95bd2b643d5369d929438da96d5bbb7cc2c918f74" +checksum = "bc9ce40b16101cb6ea63d3e221567affd1c3a9205f95d7bc574941a10636b632" dependencies = [ "anyhow", "cargo_toml", @@ -3936,15 +4132,14 @@ dependencies = [ "serde_json", "tauri-utils", "tauri-winres", - "toml 0.9.12+spec-1.1.0", "walkdir", ] [[package]] name = "tauri-codegen" -version = "2.5.4" +version = "2.6.3" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "aac423e5859d9f9ccdd32e3cf6a5866a15bedbf25aa6630bcb2acde9468f6ae3" +checksum = "08279169ff42f8fc45a1dbc9dcae888893ba95288142e5880c59b93a26d2cfc5" dependencies = [ "base64 0.22.1", "brotli", @@ -3969,9 +4164,9 @@ dependencies = [ [[package]] name = "tauri-macros" -version = "2.5.4" +version = "2.6.3" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "1b6a1bd2861ff0c8766b1d38b32a6a410f6dc6532d4ef534c47cfb2236092f59" +checksum = "e8b394794f399a421811d06966343e7933fcae92d59f5180b9388d1174497a45" dependencies = [ "heck 0.5.0", "proc-macro2", @@ -4116,9 +4311,9 @@ dependencies = [ [[package]] name = "tauri-runtime" -version = "2.10.0" +version = "2.11.3" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "b885ffeac82b00f1f6fd292b6e5aabfa7435d537cef57d11e38a489956535651" +checksum = "b0b4bc95aed361b0019067d189a1174a603d460d0f6c72606512d59fc9c12ec8" dependencies = [ "cookie", "dpi", @@ -4141,9 +4336,9 @@ dependencies = [ [[package]] name = "tauri-runtime-wry" -version = "2.10.0" +version = "2.11.4" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "5204682391625e867d16584fedc83fc292fb998814c9f7918605c789cd876314" +checksum = "4e6fac707727b7a2f48e4ded90976324267371073edbb415ffb73bb0458d203f" dependencies = [ "gtk", "http", @@ -4151,7 +4346,6 @@ dependencies = [ "log", "objc2", "objc2-app-kit", - "objc2-foundation", "once_cell", "percent-encoding", "raw-window-handle", @@ -4168,24 +4362,26 @@ dependencies = [ [[package]] name = "tauri-utils" -version = "2.8.2" +version = "2.9.3" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "fcd169fccdff05eff2c1033210b9b94acd07a47e6fa9a3431cf09cfd4f01c87e" +checksum = "3e176a18e67764923c4f1ce66f25ae4abe5f688384d5eb1a0fa6c77f3d90f887" dependencies = [ "anyhow", "brotli", "cargo_metadata", "ctor", + "dom_query", "dunce", "glob", - "html5ever", + "html5ever 0.29.1", "http", "infer", "json-patch", "kuchikiki", "log", "memchr", - "phf 0.11.3", + "phf 0.13.1", + "plist", "proc-macro2", "quote", "regex", @@ -4239,6 +4435,15 @@ dependencies = [ "utf-8", ] +[[package]] +name = "tendril" +version = "0.5.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "5fed54709c5b3a53d09bb1c113ea4f5ceafd1e772ddcb0030a82e1d56c087b08" +dependencies = [ + "new_debug_unreachable", +] + [[package]] name = "thiserror" version = "1.0.69" @@ -4531,9 +4736,9 @@ dependencies = [ [[package]] name = "tray-icon" -version = "0.21.3" +version = "0.23.1" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "a5e85aa143ceb072062fc4d6356c1b520a51d636e7bc8e77ec94be3608e5e80c" +checksum = "15edbb0d80583e85ee8df283410038e17314df5cba30da2087a54a85216c0773" dependencies = [ "crossbeam-channel", "dirs", @@ -4545,10 +4750,10 @@ dependencies = [ "objc2-core-graphics", "objc2-foundation", "once_cell", - "png 0.17.16", + "png 0.18.1", "serde", "thiserror 2.0.18", - "windows-sys 0.60.2", + "windows-sys 0.61.2", ] [[package]] @@ -4891,6 +5096,18 @@ dependencies = [ "wasm-bindgen", ] +[[package]] +name = "web_atoms" +version = "0.2.6" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "ba8b815c1b593dc0baf78dd0f4fc8fdb2de53198fb1163738093e9a311c33fb3" +dependencies = [ + "phf 0.13.1", + "phf_codegen 0.13.1", + "string_cache 0.9.0", + "string_cache_codegen 0.6.1", +] + [[package]] name = "webkit2gtk" version = "2.0.2" @@ -5538,24 +5755,23 @@ checksum = "9edde0db4769d2dc68579893f2306b26c6ecfbe0ef499b013d731b7b9247e0b9" [[package]] name = "wry" -version = "0.54.2" +version = "0.55.1" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "bb26159b420aa77684589a744ae9a9461a95395b848764ad12290a14d960a11a" +checksum = "186f9871daa55fd9c016578b810d149de58367113db7fb72b462d2323ce19514" dependencies = [ "base64 0.22.1", "block2", "cookie", "crossbeam-channel", "dirs", + "dom_query", "dpi", "dunce", "gdkx11", "gtk", - "html5ever", "http", "javascriptcore-rs", "jni", - "kuchikiki", "libc", "ndk", "objc2", diff --git a/desktop/src-tauri/Cargo.toml b/desktop/src-tauri/Cargo.toml index c2ea4cc37e..32702a0107 100644 --- a/desktop/src-tauri/Cargo.toml +++ b/desktop/src-tauri/Cargo.toml @@ -13,14 +13,15 @@ crate-type = ["staticlib", "cdylib", "rlib"] tauri-build = { version = "2", features = [] } [dependencies] -tauri = { version = "2", features = [] } -tauri-plugin-shell = "2" -tauri-plugin-opener = "2" -tauri-plugin-updater = "2" -tauri-plugin-dialog = "2" +tauri = { version = "2.11.1", features = [] } +tauri-plugin-shell = "2.3.1" +tauri-plugin-opener = "2.3.1" +tauri-plugin-updater = "2.3.1" +tauri-plugin-dialog = "2.3.1" +serde = { version = "1", features = ["derive"] } serde_json = "1" tempfile = "3" tokio = { version = "1", features = ["time", "sync"] } [target.'cfg(target_os = "macos")'.dependencies] -tauri = { version = "2", features = ["image-png", "tray-icon"] } +tauri = { version = "2.11.1", features = ["image-png", "tray-icon"] } diff --git a/desktop/src-tauri/build.rs b/desktop/src-tauri/build.rs index d860e1e6a7..2e427d37bd 100644 --- a/desktop/src-tauri/build.rs +++ b/desktop/src-tauri/build.rs @@ -1,3 +1,11 @@ fn main() { - tauri_build::build() + tauri_build::try_build(tauri_build::Attributes::new().app_manifest( + tauri_build::AppManifest::new().commands(&[ + "claude_auth_start", + "claude_auth_disconnect", + "claude_auth_status", + "claude_auth_fetch_result", + ]), + )) + .expect("failed to build Tauri application manifest"); } diff --git a/desktop/src-tauri/capabilities/claude-auth.json b/desktop/src-tauri/capabilities/claude-auth.json new file mode 100644 index 0000000000..8ad65b5e15 --- /dev/null +++ b/desktop/src-tauri/capabilities/claude-auth.json @@ -0,0 +1,15 @@ +{ + "$schema": "../gen/schemas/desktop-schema.json", + "identifier": "claude-auth", + "description": "Claude authentication webview may return an authenticated fetch result only.", + "remote": { + "urls": [ + "https://claude.ai/*", + "https://*.claude.ai/*" + ] + }, + "windows": ["claude-auth"], + "permissions": [ + "allow-claude-auth-fetch-result" + ] +} diff --git a/desktop/src-tauri/capabilities/default.json b/desktop/src-tauri/capabilities/default.json index 6e5f25e781..26509425bc 100644 --- a/desktop/src-tauri/capabilities/default.json +++ b/desktop/src-tauri/capabilities/default.json @@ -10,6 +10,9 @@ "windows": ["main"], "permissions": [ "core:default", - "core:webview:allow-set-webview-zoom" + "core:webview:allow-set-webview-zoom", + "allow-claude-auth-start", + "allow-claude-auth-disconnect", + "allow-claude-auth-status" ] } diff --git a/desktop/src-tauri/src/lib.rs b/desktop/src-tauri/src/lib.rs index 594ab8d700..c267a97e0b 100644 --- a/desktop/src-tauri/src/lib.rs +++ b/desktop/src-tauri/src/lib.rs @@ -20,7 +20,10 @@ use tauri::menu::{MenuBuilder, MenuItemBuilder, SubmenuBuilder}; use tauri::plugin::Builder as PluginBuilder; #[cfg(target_os = "macos")] use tauri::tray::TrayIconBuilder; -use tauri::{App, AppHandle, Emitter, Manager, RunEvent, Url, WebviewWindow}; +use tauri::{ + App, AppHandle, Emitter, Manager, RunEvent, State, Url, WebviewUrl, WebviewWindow, + WebviewWindowBuilder, +}; use tauri_plugin_dialog::{DialogExt, MessageDialogButtons}; use tauri_plugin_opener::OpenerExt; use tauri_plugin_shell::process::{CommandChild, CommandEvent}; @@ -54,6 +57,10 @@ const CHECK_UPDATES_MENU_ID: &str = "check_updates"; const OPEN_LOGS_FOLDER_MENU_ID: &str = "open_logs_folder"; const SHOW_MAIN_WINDOW_MENU_ID: &str = "show_main_window"; const QUIT_FROM_STATUS_ITEM_MENU_ID: &str = "quit_from_status_item"; +const CLAUDE_AUTH_WINDOW_LABEL: &str = "claude-auth"; +const CLAUDE_AUTH_URL: &str = "https://claude.ai/login?return_url=%2Fnew"; +const CLAUDE_BROWSER_FETCH_TIMEOUT: Duration = Duration::from_secs(45); +const CLAUDE_BROWSER_RESPONSE_MAX_BYTES: usize = 32 * 1024 * 1024; // Delay after navigating to the backend before probing whether the // Linux WebKitGTK web content process is actually alive. Gives the // process time to spawn so we don't false-positive on slow startup. @@ -67,6 +74,9 @@ type CommandRx = Receiver; struct SidecarState { child: Mutex>, backend_port: Mutex>, + // Resolved once from the exact environment supplied to the Go sidecar. + // It is held in memory only and cleared when that sidecar terminates. + loopback_auth_token: Mutex>, active_generation: Mutex>, stopping_generation: Mutex>, restart_after_stop_timeout_generation: Mutex>, @@ -82,6 +92,56 @@ struct SidecarProcess { generation: u64, } +/// Deliberately holds only non-secret state. Claude session cookies stay in the +/// isolated WKWebView profile and are never exposed to the frontend or logs. +#[derive(Default)] +struct ClaudeAuthState { + connected_this_launch: AtomicBool, + auth_watcher_active: AtomicBool, + transport_worker_active: AtomicBool, + next_browser_request: AtomicU64, + pending_browser_request: Mutex>, +} + +struct ClaudeBrowserRequest { + id: u64, + response: SyncSender, +} + +struct ClaudeBrowserResponse { + status: u16, + body: String, + error: Option, + retry_after: Option, +} + +#[derive(serde::Deserialize)] +#[serde(rename_all = "camelCase")] +struct ClaudeTransportRequest { + id: String, + provider: String, + operation: String, + params: serde_json::Value, + lease: String, +} + +#[derive(serde::Deserialize)] +#[serde(rename_all = "camelCase")] +struct ClaudeBrowserFetchResult { + request_id: u64, + status: u16, + body: String, + error: Option, + retry_after: Option, +} + +#[derive(serde::Serialize)] +#[serde(rename_all = "camelCase")] +struct ClaudeAuthStatus { + connected: bool, + message: String, +} + #[derive(Clone, Copy, Debug, PartialEq, Eq)] enum DesktopMenuAction { About, @@ -153,6 +213,13 @@ pub fn run() { .plugin(tauri_plugin_dialog::init()) .plugin(init_navigation_guard_plugin()) .manage(SidecarState::default()) + .manage(ClaudeAuthState::default()) + .invoke_handler(tauri::generate_handler![ + claude_auth_start, + claude_auth_disconnect, + claude_auth_status, + claude_auth_fetch_result, + ]) .setup(|app| { if let Err(err) = setup_menu(app) { eprintln!("[agentsview] failed to set up desktop menu: {err}"); @@ -177,6 +244,10 @@ pub fn run() { err.to_string().as_str(), ); } else { + // Restore the isolated Claude profile in a hidden window so + // scheduled syncs can run after an app restart. The watcher + // makes it visible only when interactive sign-in is needed. + start_claude_auth_background(app.handle().clone()); schedule_auto_update_check(app.handle().clone()); } } @@ -442,6 +513,12 @@ fn combined_preflight_output(stdout: &str, stderr: &str) -> Option { fn init_navigation_guard_plugin() -> tauri::plugin::TauriPlugin { PluginBuilder::new("navigation-guard") .on_navigation(|webview, url| { + // The auth window may visit Claude plus selected identity providers. + // Only Claude origins receive the fetch-result command capability; + // identity-provider documents have no native IPC permissions. + if webview.label() == CLAUDE_AUTH_WINDOW_LABEL { + return is_allowed_claude_auth_navigation(url); + } let backend_port = webview .app_handle() .try_state::() @@ -491,6 +568,514 @@ fn is_allowed_navigation_url(url: &Url, backend_port: Option) -> bool { false } +fn is_allowed_claude_auth_navigation(url: &Url) -> bool { + if url.scheme() != "https" { + return false; + } + let Some(host) = url.host_str() else { + return false; + }; + host == "claude.ai" + || host.ends_with(".claude.ai") + || matches!( + host, + "accounts.google.com" | "login.microsoftonline.com" | "appleid.apple.com" + ) + || host.ends_with(".okta.com") + || host.ends_with(".auth0.com") +} + +fn claude_auth_profile_dir(handle: &AppHandle) -> Result { + let directory = handle + .path() + .app_data_dir() + .map_err(|err| format!("could not resolve the Claude login profile directory: {err}"))? + .join("cloud-auth") + .join("claude-ai"); + fs::create_dir_all(&directory) + .map_err(|err| format!("could not create the Claude login profile directory: {err}"))?; + #[cfg(unix)] + fs::set_permissions(&directory, fs::Permissions::from_mode(0o700)) + .map_err(|err| format!("could not secure the Claude login profile directory: {err}"))?; + Ok(directory) +} + +fn create_claude_auth_window(handle: &AppHandle, visible: bool) -> Result { + let profile_dir = claude_auth_profile_dir(handle)?; + let url = Url::parse(CLAUDE_AUTH_URL).map_err(|err| err.to_string())?; + WebviewWindowBuilder::new(handle, CLAUDE_AUTH_WINDOW_LABEL, WebviewUrl::External(url)) + .title("Connect Claude.ai to AgentsView") + .inner_size(1100.0, 800.0) + .min_inner_size(800.0, 600.0) + .visible(visible) + .data_directory(profile_dir) + .build() + .map_err(|err| format!("could not open the Claude sign-in window: {err}")) +} + +fn start_claude_auth_background(handle: AppHandle) { + if handle + .get_webview_window(CLAUDE_AUTH_WINDOW_LABEL) + .is_some() + { + return; + } + if create_claude_auth_window(&handle, false).is_ok() { + start_claude_auth_watcher(handle); + } +} + +#[tauri::command] +fn claude_auth_start(handle: AppHandle) -> Result { + if let Some(window) = handle.get_webview_window(CLAUDE_AUTH_WINDOW_LABEL) { + window.show().map_err(|err| err.to_string())?; + window.set_focus().map_err(|err| err.to_string())?; + if !handle + .state::() + .connected_this_launch + .load(Ordering::SeqCst) + { + start_claude_auth_watcher(handle.clone()); + } + return Ok(ClaudeAuthStatus { + connected: handle.state::().connected_this_launch.load(Ordering::SeqCst), + message: "Claude sign-in is already open. AgentsView will verify the session automatically once sign-in completes.".into(), + }); + } + + create_claude_auth_window(&handle, true)?; + start_claude_auth_watcher(handle.clone()); + Ok(ClaudeAuthStatus { + connected: false, + message: "Sign in to Claude in the separate window. AgentsView will verify the session automatically.".into(), + }) +} + +fn loopback_auth_token(handle: &AppHandle) -> Option { + handle + .state::() + .loopback_auth_token + .lock() + .ok() + .and_then(|token| token.clone()) +} + +// Resolve from the same merged environment passed to Go, rather than Tauri's +// app-data directory. This covers AGENTSVIEW_AUTH_TOKEN and +// AGENTSVIEW_DATA_DIR overrides; the resulting value is retained only for the +// lifetime of the matching sidecar generation. +fn resolve_sidecar_loopback_auth_token() -> Option { + let environment = sidecar_env(); + let lookup = |name: &str| { + environment + .iter() + .find_map(|(key, value)| (key == name).then(|| value.to_string_lossy().into_owned())) + .filter(|value| !value.is_empty()) + }; + if let Some(token) = lookup("AGENTSVIEW_AUTH_TOKEN") { + return Some(token); + } + let data_dir = lookup("AGENTSVIEW_DATA_DIR") + .or_else(|| lookup("AGENT_VIEWER_DATA_DIR")) + .map(PathBuf::from) + .or_else(|| resolve_home_dir().map(|home| home.join(".agentsview")))?; + let config = fs::read_to_string(data_dir.join("config.toml")).ok()?; + config.lines().find_map(|line| { + let (key, value) = line.split_once('=')?; + (key.trim() == "auth_token") + .then(|| value.trim().trim_matches('"').to_string()) + .filter(|token| !token.is_empty()) + }) +} + +struct LoopbackResponse { + body: String, +} + +fn claude_loopback_post( + handle: &AppHandle, + path: &str, + payload: serde_json::Value, +) -> Result { + let port = handle + .state::() + .backend_port + .lock() + .ok() + .and_then(|v| *v) + .ok_or_else(|| "AgentsView backend is not running".to_string())?; + let body = serde_json::to_vec(&payload).map_err(|err| err.to_string())?; + let mut stream = TcpStream::connect((HOST, port)) + .map_err(|err| format!("connect local AgentsView backend: {err}"))?; + let request = loopback_http_request( + path, + port, + body.len(), + loopback_auth_token(handle).as_deref(), + ); + stream + .write_all(request.as_bytes()) + .and_then(|_| stream.write_all(&body)) + .map_err(|err| err.to_string())?; + let mut response = String::new(); + stream + .read_to_string(&mut response) + .map_err(|err| err.to_string())?; + parse_loopback_response(&response) +} + +fn parse_loopback_response(response: &str) -> Result { + let (head, body) = response + .split_once("\r\n\r\n") + .ok_or_else(|| "invalid local backend response".to_string())?; + let status = head + .lines() + .next() + .and_then(|line| line.split_whitespace().nth(1)) + .and_then(|status| status.parse::().ok()) + .ok_or_else(|| "invalid local backend status".to_string())?; + if !(200..300).contains(&status) { + return Err(format!( + "local backend request failed: {}", + head.lines().next().unwrap_or("unknown") + )); + } + Ok(LoopbackResponse { + body: body.to_string(), + }) +} + +fn loopback_http_request(path: &str, port: u16, body_len: usize, token: Option<&str>) -> String { + let authorization = token + .map(|value| format!("Authorization: Bearer {value}\r\n")) + .unwrap_or_default(); + format!("POST {path} HTTP/1.1\r\nHost: {HOST}:{port}\r\nOrigin: http://{HOST}:{port}\r\n{authorization}Content-Type: application/json\r\nContent-Length: {body_len}\r\nConnection: close\r\n\r\n") +} + +// start_claude_transport_worker is deliberately transport-only: it claims +// typed work from the Go sync service and returns one browser response. It +// contains no pagination, retry, scheduling, cache, or import policy. +fn start_claude_transport_worker(handle: AppHandle) { + let state = handle.state::(); + if state.transport_worker_active.swap(true, Ordering::SeqCst) { + return; + } + thread::spawn(move || loop { + if !handle + .state::() + .connected_this_launch + .load(Ordering::SeqCst) + { + thread::sleep(Duration::from_secs(1)); + continue; + } + let claimed = match claude_loopback_post( + &handle, + "/api/v1/cloud/transport/claim", + serde_json::json!({}), + ) { + Ok(response) => response, + Err(_) => { + thread::sleep(Duration::from_millis(250)); + continue; + } + }; + let Ok(request) = serde_json::from_str::(&claimed.body) else { + // A successful claim with an unusable body cannot represent work. + // Back off to avoid a tight local loop while the daemon is restarting. + thread::sleep(Duration::from_millis(250)); + continue; + }; + if request.id.is_empty() { + continue; + } + let response = claude_execute_transport_request(&handle, &request); + let payload = match response { + Ok(response) => { + serde_json::json!({"id":request.id,"lease":request.lease,"status":response.status,"body":serde_json::from_str::(&response.body).unwrap_or(serde_json::Value::Null),"retry_after":response.retry_after,"error":response.error}) + } + Err(error) => { + serde_json::json!({"id":request.id,"lease":request.lease,"status":0,"body":null,"error":error}) + } + }; + if let Err(err) = claude_loopback_post(&handle, "/api/v1/cloud/transport/result", payload) { + // Never include browser response bodies or credentials in diagnostics. + // This is intentionally visible: a failed completion otherwise leaves + // the Go broker waiting for its lease to expire with no useful signal. + eprintln!("[agentsview] Claude transport result delivery failed: {err}"); + } + }); +} + +fn claude_execute_transport_request( + handle: &AppHandle, + request: &ClaudeTransportRequest, +) -> Result { + let path = validate_claude_transport_request(handle, request)?; + claude_browser_fetch(handle, &path) +} + +fn validate_claude_transport_request( + handle: &AppHandle, + request: &ClaudeTransportRequest, +) -> Result { + if request.provider != "claude-ai" { + return Err("unsupported cloud provider".into()); + } + let organization = claude_session_organization(handle)?; + match request.operation.as_str() { + "list_conversations" => { + let offset = request + .params + .get("offset") + .and_then(serde_json::Value::as_u64); + let limit = request + .params + .get("limit") + .and_then(serde_json::Value::as_u64); + if offset.is_none() || !matches!(limit, Some(1..=100)) { + return Err("invalid Claude list request".into()); + } + Ok(format!("/api/organizations/{organization}/chat_conversations_v2?limit={}&offset={}&starred=false&consistency=eventual", limit.unwrap(), offset.unwrap())) + } + "get_conversation" => { + let id = request + .params + .get("conversation_id") + .and_then(serde_json::Value::as_str); + if !id.is_some_and(valid_claude_identifier) { + return Err("invalid Claude conversation id".into()); + } + Ok(format!("/api/organizations/{organization}/chat_conversations/{}?tree=True&rendering_mode=messages&consistency=strong", id.unwrap())) + } + _ => Err("unsupported Claude transport operation".into()), + } +} + +fn valid_claude_identifier(value: &str) -> bool { + !value.is_empty() + && value.len() <= 128 + && value + .chars() + .all(|character| character.is_ascii_alphanumeric() || character == '-') +} + +// Read only session markers required to verify sign-in and scope the typed API +// operation. They stay inside this process and are never persisted or logged. +fn claude_session_organization(handle: &AppHandle) -> Result { + let window = handle + .get_webview_window(CLAUDE_AUTH_WINDOW_LABEL) + .ok_or_else(|| "Claude sign-in window was closed.".to_string())?; + let url = Url::parse("https://claude.ai/").map_err(|err| err.to_string())?; + let cookies = window + .cookies_for_url(url) + .map_err(|err| format!("could not inspect Claude sign-in state: {err}"))?; + let signed_in = cookies.iter().any(|cookie| cookie.name() == "sessionKey"); + let organization = cookies + .iter() + .find(|cookie| cookie.name() == "lastActiveOrg") + .map(|cookie| cookie.value().to_string()) + .filter(|value| valid_claude_identifier(value)); + if !signed_in || organization.is_none() { + return Err("Claude is not signed in yet.".into()); + } + Ok(organization.expect("checked above")) +} + +fn claude_browser_fetch(handle: &AppHandle, path: &str) -> Result { + if !path.starts_with("/api/organizations/") || path.contains('\n') || path.contains('\r') { + return Err("unsupported Claude API path".into()); + } + let window = handle + .get_webview_window(CLAUDE_AUTH_WINDOW_LABEL) + .ok_or_else(|| { + "Claude browser session is not available. Reconnect Claude and try again.".to_string() + })?; + let state = handle.state::(); + let request_id = state.next_browser_request.fetch_add(1, Ordering::SeqCst); + let (sender, receiver) = sync_channel(1); + { + let mut pending = state + .pending_browser_request + .lock() + .map_err(|_| "Claude browser request lock failed")?; + if pending.is_some() { + return Err("another Claude browser request is already running".into()); + } + *pending = Some(ClaudeBrowserRequest { + id: request_id, + response: sender, + }); + } + + let url = format!("https://claude.ai{path}"); + let url = serde_json::to_string(&url).map_err(|err| err.to_string())?; + let script = format!( + r#"(async () => {{ + try {{ + const response = await fetch({url}, {{ credentials: "include" }}); + const body = await response.text(); + const retryAfter = response.headers.get("retry-after") ?? undefined; + await window.__TAURI__.core.invoke("claude_auth_fetch_result", {{ payload: {{ requestId: {request_id}, status: response.status, body, retryAfter }} }}); + }} catch (error) {{ + await window.__TAURI__.core.invoke("claude_auth_fetch_result", {{ payload: {{ requestId: {request_id}, status: 0, body: "", error: String(error), retryAfter: undefined }} }}); + }} + }})()"#, + ); + if let Err(err) = window.eval(&script) { + let mut pending = state + .pending_browser_request + .lock() + .map_err(|_| "Claude browser request lock failed")?; + *pending = None; + return Err(format!("could not start Claude browser request: {err}")); + } + match receiver.recv_timeout(CLAUDE_BROWSER_FETCH_TIMEOUT) { + Ok(response) => Ok(response), + Err(_) => { + let mut pending = state + .pending_browser_request + .lock() + .map_err(|_| "Claude browser request lock failed")?; + *pending = None; + Err("Claude browser request timed out".into()) + } + } +} + +#[tauri::command] +fn claude_auth_fetch_result( + window: WebviewWindow, + state: State<'_, ClaudeAuthState>, + payload: ClaudeBrowserFetchResult, +) -> Result<(), String> { + if window.label() != CLAUDE_AUTH_WINDOW_LABEL { + return Err("Claude browser response came from an unexpected window".into()); + } + let origin = window.url().map_err(|err| err.to_string())?; + if origin.scheme() != "https" + || !origin + .host_str() + .is_some_and(|host| host == "claude.ai" || host.ends_with(".claude.ai")) + { + return Err("Claude browser response came from an unexpected origin".into()); + } + if payload.body.len() > CLAUDE_BROWSER_RESPONSE_MAX_BYTES { + let mut pending = state + .pending_browser_request + .lock() + .map_err(|_| "Claude browser request lock failed")?; + if pending + .as_ref() + .is_some_and(|request| request.id == payload.request_id) + { + let request = pending.take().expect("checked above"); + let _ = request.response.send(ClaudeBrowserResponse { + status: 0, + body: String::new(), + error: Some("Claude browser response exceeded the 32 MiB safety limit".into()), + retry_after: None, + }); + } + return Err("Claude browser response was too large".into()); + } + let mut pending = state + .pending_browser_request + .lock() + .map_err(|_| "Claude browser request lock failed")?; + let Some(request) = pending.take() else { + return Err("no Claude browser request is pending".into()); + }; + if request.id != payload.request_id { + *pending = Some(request); + return Err("Claude browser response did not match the pending request".into()); + } + request + .response + .send(ClaudeBrowserResponse { + status: payload.status, + body: payload.body, + error: payload.error, + retry_after: payload.retry_after, + }) + .map_err(|_| "Claude browser request was cancelled".into()) +} + +fn start_claude_auth_watcher(handle: AppHandle) { + let state = handle.state::(); + if state.auth_watcher_active.swap(true, Ordering::SeqCst) { + return; + } + thread::spawn(move || { + // Poll only while this dedicated sign-in window exists, and stop after + // ten minutes. Reopening an unauthenticated window starts a new watcher. + for _ in 0..600 { + thread::sleep(Duration::from_secs(1)); + let authenticated = claude_session_organization(&handle).is_ok(); + if !authenticated { + continue; + } + handle + .state::() + .connected_this_launch + .store(true, Ordering::SeqCst); + if let Some(window) = handle.get_webview_window(CLAUDE_AUTH_WINDOW_LABEL) { + // Keep the isolated webview alive, but out of the user's way: + // the Go sync service leases typed requests to this transport. + let _ = window.hide(); + } + start_claude_transport_worker(handle.clone()); + break; + } + handle + .state::() + .auth_watcher_active + .store(false, Ordering::SeqCst); + }); +} + +#[tauri::command] +fn claude_auth_disconnect(handle: AppHandle) -> Result { + if let Some(window) = handle.get_webview_window(CLAUDE_AUTH_WINDOW_LABEL) { + // On macOS, deleting the WebKit profile directory alone can leave the + // persistent WKWebsiteDataStore (and its Claude session) intact. + window + .clear_all_browsing_data() + .map_err(|err| format!("could not clear the Claude browser session: {err}"))?; + window.close().map_err(|err| err.to_string())?; + } + let profile_dir = claude_auth_profile_dir(&handle)?; + if profile_dir.exists() { + fs::remove_dir_all(&profile_dir) + .map_err(|err| format!("could not remove the Claude login profile: {err}"))?; + } + handle + .state::() + .connected_this_launch + .store(false, Ordering::SeqCst); + Ok(ClaudeAuthStatus { + connected: false, + message: "Claude has been disconnected and its isolated browser session was removed." + .into(), + }) +} + +#[tauri::command] +fn claude_auth_status(handle: AppHandle) -> ClaudeAuthStatus { + let connected = handle + .state::() + .connected_this_launch + .load(Ordering::SeqCst); + ClaudeAuthStatus { + connected, + message: if connected { + "Claude session is held in an isolated browser profile on this device.".into() + } else { + "Not connected.".into() + }, + } +} + fn is_allowed_external_open_url(url: &Url) -> bool { matches!( url.scheme(), @@ -878,6 +1463,9 @@ fn save_sidecar(app: &AppHandle, child: CommandChild) -> Result { if let Ok(mut active_generation) = state.active_generation.lock() { *active_generation = Some(generation); } + if let Ok(mut token) = state.loopback_auth_token.lock() { + *token = resolve_sidecar_loopback_auth_token(); + } if let Ok(mut stopping_generation) = state.stopping_generation.lock() { *stopping_generation = None; } @@ -903,6 +1491,12 @@ fn set_sidecar_port(state: &SidecarState, port: Option) { } } +fn clear_sidecar_loopback_auth_token(state: &SidecarState) { + if let Ok(mut token) = state.loopback_auth_token.lock() { + *token = None; + } +} + fn handle_sidecar_terminated( state: &SidecarState, startup_handled: &AtomicBool, @@ -910,6 +1504,7 @@ fn handle_sidecar_terminated( ) -> bool { if mark_sidecar_inactive_if_current(state, generation) { set_sidecar_port(state, None); + clear_sidecar_loopback_auth_token(state); } clear_sidecar_child_if_current(state, generation); clear_stopping_generation_if_current(state, generation); @@ -3039,6 +3634,40 @@ mod tests { use std::time::{SystemTime, UNIX_EPOCH}; use tempfile::tempdir; + #[test] + fn claude_loopback_request_has_pinned_origin_and_bearer_token() { + let request = + loopback_http_request("/api/v1/cloud/transport/claim", 18080, 2, Some("secret")); + assert!(request.contains("Origin: http://127.0.0.1:18080\r\n")); + assert!(request.contains("Authorization: Bearer secret\r\n")); + assert!(!request.contains("token=")); + } + + #[test] + fn loopback_response_accepts_empty_204_completion() { + let response = + parse_loopback_response("HTTP/1.1 204 No Content\r\nContent-Length: 0\r\n\r\n") + .expect("204 response must complete transport work"); + assert!(response.body.is_empty()); + } + + #[test] + fn loopback_response_rejects_non_success_status() { + assert!(parse_loopback_response("HTTP/1.1 401 Unauthorized\r\n\r\n{}").is_err()); + } + + #[test] + fn claude_browser_response_limit_allows_large_conversation_payloads() { + assert!(CLAUDE_BROWSER_RESPONSE_MAX_BYTES >= 32 * 1024 * 1024); + } + + #[test] + fn claude_identifier_validation_is_allow_listed() { + assert!(valid_claude_identifier("abc-123")); + assert!(!valid_claude_identifier("../secret")); + assert!(!valid_claude_identifier("abc\n123")); + } + #[test] fn sidecar_args_use_cobra_long_flags() { assert_eq!( diff --git a/desktop/src-tauri/tauri.conf.json b/desktop/src-tauri/tauri.conf.json index c0d6d33281..9db4b2419f 100644 --- a/desktop/src-tauri/tauri.conf.json +++ b/desktop/src-tauri/tauri.conf.json @@ -1,7 +1,7 @@ { "$schema": "https://schema.tauri.app/config/2", "productName": "AgentsView", - "version": "0.12.1", + "version": "0.40.1-dev.16", "identifier": "io.agentsview.desktop", "build": { "frontendDist": "../ui" diff --git a/docs/claude-ai.md b/docs/claude-ai.md new file mode 100644 index 0000000000..d86e6e5b89 --- /dev/null +++ b/docs/claude-ai.md @@ -0,0 +1,7 @@ +# Claude.ai sync + +Claude.ai imports use an isolated authenticated desktop WKWebView. Session cookies remain in that browser profile and are never sent to the AgentsView daemon, cache, database, or scheduler configuration. + +The Go server owns the sync job: pagination, retries, cache markers, repair, cancellation, scheduling and SQLite import. Tauri is only an authenticated transport adapter: it receives an allow-listed typed request from the local server and returns the Claude response. This permits another desktop transport to serve the same job model later without moving archive logic out of the generic web application. + +Automatic sync is optional and runs only while AgentsView is running. The Go scheduler stores its credential-free configuration under `cloud-cache/claude-ai/schedule.json`; it waits for an authenticated desktop transport rather than reading browser credentials. The cache records durable import progress and summary markers, so a later scan safely skips unchanged conversation details. **Repair import** explicitly refetches every conversation, replaces cached content even when timestamps match, and clears only matching Claude.ai permanent-deletion markers. diff --git a/frontend/src/lib/api/client.ts b/frontend/src/lib/api/client.ts index 5574dade2c..7a777e611c 100644 --- a/frontend/src/lib/api/client.ts +++ b/frontend/src/lib/api/client.ts @@ -582,6 +582,63 @@ async function readImportSSE( return result; } +export interface ClaudeSyncStatus { + id: string; + status: "idle" | "running" | "cancelling" | "cancelled" | "completed" | "failed"; + mode: "incremental" | "repair"; + scanned: number; + changed: number; + fetched: number; + imported: number; + updated: number; + skipped: number; + failed: number; + error?: string; +} + +export async function startClaudeAISync(mode: "incremental" | "repair"): Promise { + const res = await fetch(`${getBase()}/cloud/claude-ai/sync`, authHeaders({ + method: "POST", + headers: { "Content-Type": "application/json" }, + body: JSON.stringify({ mode }), + })); + if (!res.ok) throw new Error(`Claude sync could not start (${res.status})`); + return res.json(); +} + +export async function getClaudeAISyncStatus(): Promise { + const res = await fetch(`${getBase()}/cloud/claude-ai/status`, authHeaders()); + if (!res.ok) throw new Error(`Claude sync status failed (${res.status})`); + return res.json(); +} + +export async function cancelClaudeAISync(): Promise { + const res = await fetch(`${getBase()}/cloud/claude-ai/cancel`, authHeaders({ method: "POST" })); + if (!res.ok) throw new Error(`Claude sync cancellation failed (${res.status})`); + return res.json(); +} + +export interface ClaudeScheduleConfig { + enabled: boolean; + interval_minutes: number; +} + +export async function getClaudeAISchedule(): Promise { + const res = await fetch(`${getBase()}/cloud/claude-ai/schedule`, authHeaders()); + if (!res.ok) throw new Error(`Claude schedule lookup failed (${res.status})`); + return res.json(); +} + +export async function configureClaudeAISchedule(config: ClaudeScheduleConfig): Promise { + const res = await fetch(`${getBase()}/cloud/claude-ai/schedule`, authHeaders({ + method: "POST", + headers: { "Content-Type": "application/json" }, + body: JSON.stringify(config), + })); + if (!res.ok) throw new Error(`Claude schedule update failed (${res.status})`); + return res.json(); +} + export async function importClaudeAI( file: File, cb?: ImportCallbacks, diff --git a/frontend/src/lib/components/settings/ClaudeAiSettings.svelte b/frontend/src/lib/components/settings/ClaudeAiSettings.svelte new file mode 100644 index 0000000000..f86ec52a58 --- /dev/null +++ b/frontend/src/lib/components/settings/ClaudeAiSettings.svelte @@ -0,0 +1,326 @@ + + +
+
+
+ + {status.connected ? "Connected" : "Not connected"} +
+

+ Sign in through an isolated Claude.ai window. Its browser session stays in + that profile and is never shown here. +

+ {#if !isDesktop} +

Claude connection is available in the AgentsView desktop app.

+ {:else if !canUseClaudeBrowser} +

Claude connection is unavailable while viewing a remote AgentsView server.

+ {:else} +

{status.message}

+ {/if} + {#if error} + + {/if} + {#if syncActive} +
+ +
+ {sync.status === "cancelling" ? "Cancelling Claude sync…" : "Syncing Claude conversations…"} + + Scanned {sync.scanned} · changed {sync.changed} · fetched {sync.fetched} · imported {sync.imported + sync.updated} · skipped {sync.skipped} + + You can leave this page. The sync continues while AgentsView is running. +
+
+ {:else if sync.status === "completed" || sync.status === "cancelled" || sync.status === "failed"} +

+ Claude sync {sync.status}: scanned {sync.scanned}, changed {sync.changed}, imported {sync.imported + sync.updated}, skipped {sync.skipped}, failed {sync.failed}. +

+ {/if} +
+ + {#if canUseClaudeBrowser} +
+ + {#if status.connected} + + + {#if busy || syncActive} + + {:else if error} + + {/if} + + {/if} +
+ + {/if} +
+ + diff --git a/frontend/src/lib/components/settings/SettingsPage.svelte b/frontend/src/lib/components/settings/SettingsPage.svelte index debca47da6..e5a3c78641 100644 --- a/frontend/src/lib/components/settings/SettingsPage.svelte +++ b/frontend/src/lib/components/settings/SettingsPage.svelte @@ -20,6 +20,7 @@ import TerminalSettings from "./TerminalSettings.svelte"; import EmbeddingsSettings from "./EmbeddingsSettings.svelte"; import GithubSettings from "./GithubSettings.svelte"; + import ClaudeAiSettings from "./ClaudeAiSettings.svelte"; import LanguageSettings from "./LanguageSettings.svelte"; import RemoteSettings from "./RemoteSettings.svelte"; import { settingsPanels } from "./settingsPanels.js"; @@ -214,6 +215,8 @@ {:else if meta.id === "github"} + {:else if meta.id === "claude-ai"} + {:else if meta.id === "remote-access"} {/if} diff --git a/frontend/src/lib/components/settings/SettingsPage.test.ts b/frontend/src/lib/components/settings/SettingsPage.test.ts index 7e23c60d13..2ebbb80439 100644 --- a/frontend/src/lib/components/settings/SettingsPage.test.ts +++ b/frontend/src/lib/components/settings/SettingsPage.test.ts @@ -263,7 +263,7 @@ describe("SettingsPage", () => { restoredSearch.dispatchEvent(new Event("input", { bubbles: true })); await tick(); - expect(restoredNav.querySelectorAll("button")).toHaveLength(9); + expect(restoredNav.querySelectorAll("button")).toHaveLength(10); expect( document.body.querySelector(".settings-page")?.classList.contains( "settings-no-results", diff --git a/frontend/src/lib/components/settings/settingsPanels.ts b/frontend/src/lib/components/settings/settingsPanels.ts index 8a7d26e999..c8de51b5a7 100644 --- a/frontend/src/lib/components/settings/settingsPanels.ts +++ b/frontend/src/lib/components/settings/settingsPanels.ts @@ -9,6 +9,7 @@ export type SettingsPanelId = | "worktree-mappings" | "embeddings" | "github" + | "claude-ai" | "remote-access"; export interface SettingsPanelMeta { @@ -90,6 +91,14 @@ export function settingsPanels(): SettingsPanelMeta[] { group: connections, keywords: m.settings_search_keywords_github(), }, + { + id: "claude-ai", + label: "Claude.ai", + title: "Claude.ai", + description: "Connect a Claude.ai session from the desktop app.", + group: connections, + keywords: "claude ai login browser session conversations import", + }, { id: "remote-access", label: m.settings_nav_remote_access(), diff --git a/go.sum b/go.sum index 4297f780f7..a75e3a9566 100644 --- a/go.sum +++ b/go.sum @@ -283,10 +283,6 @@ github.com/zeebo/xxh3 v1.1.0 h1:s7DLGDK45Dyfg7++yxI0khrfwq9661w9EN78eP/UZVs= github.com/zeebo/xxh3 v1.1.0/go.mod h1:IisAie1LELR4xhVinxWS5+zf1lA4p0MW4T+w+W07F5s= go.kenn.io/docbank v0.11.0 h1:MLv0CxCWyk5kM5zX3YALquVxi7sjAjrcMD/XNQkuw6Y= go.kenn.io/docbank v0.11.0/go.mod h1:4hym+ONhsG8epaBgTXA1oC/PnF60Z3CyvcRs44qKhb8= -go.kenn.io/kit v0.11.0 h1:OdEaI8i3R7M0OTptrP2Osu+WJSg+lTClmHalRgn/T/U= -go.kenn.io/kit v0.11.0/go.mod h1:dComZhFNb4LR+Tj4ZD0slEDMkPJk0gGd8q/HEHXTGSM= -go.kenn.io/kit v0.13.0 h1:N4/KvR1xnM2o97q2CHnBWAt1uXcNqxTFGT0dbkiloEA= -go.kenn.io/kit v0.13.0/go.mod h1:dComZhFNb4LR+Tj4ZD0slEDMkPJk0gGd8q/HEHXTGSM= go.kenn.io/kit v0.13.1 h1:KQxCS2GMczrrWhZgxjDwD480hB4nXgtw7jfBuJujhWg= go.kenn.io/kit v0.13.1/go.mod h1:dComZhFNb4LR+Tj4ZD0slEDMkPJk0gGd8q/HEHXTGSM= go.opentelemetry.io/auto/sdk v1.2.1 h1:jXsnJ4Lmnqd11kwkBV2LgLoFMZKizbCi5fNZ/ipaZ64= diff --git a/internal/cloudsync/claudeai/import.go b/internal/cloudsync/claudeai/import.go new file mode 100644 index 0000000000..bfc7121913 --- /dev/null +++ b/internal/cloudsync/claudeai/import.go @@ -0,0 +1,389 @@ +package claudeai + +import ( + "encoding/json" + "fmt" + "os" + "path/filepath" + "strings" + "time" +) + +const ( + cacheManifestName = "index.json" + syncStateName = "sync-state.json" +) + +// SyncState is the durable, credential-free progress record for a browser +// import. It is deliberately separate from the cache manifest so a stopped +// scan can safely restart from offset zero and skip unchanged details. +type SyncState struct { + Status string `json:"status"` + StartedAt time.Time `json:"started_at"` + UpdatedAt time.Time `json:"updated_at"` + Scanned int `json:"scanned"` + Changed int `json:"changed"` + Fetched int `json:"fetched"` + Imported int `json:"imported"` + Skipped int `json:"skipped"` + Failed int `json:"failed"` +} + +// BrowserImportPlan tells the authenticated browser which summary details +// actually need fetching. Summary timestamps are compared only to the local +// content manifest; no credential is accepted or persisted. +type BrowserImportPlan struct { + ChangedIDs []string `json:"changed_ids"` + Unchanged int `json:"unchanged"` + State SyncState `json:"state"` +} + +func syncStatePath(root string) string { return filepath.Join(root, syncStateName) } + +func ReadSyncState(root string) (SyncState, error) { + data, err := os.ReadFile(syncStatePath(root)) + if os.IsNotExist(err) { + return SyncState{}, nil + } + if err != nil { + return SyncState{}, fmt.Errorf("read Claude sync state: %w", err) + } + var state SyncState + if err := json.Unmarshal(data, &state); err != nil { + return SyncState{}, fmt.Errorf("decode Claude sync state: %w", err) + } + return state, nil +} + +func writeSyncState(root string, state SyncState) error { + state.UpdatedAt = time.Now().UTC() + return atomicWriteJSON(syncStatePath(root), state) +} + +// PlanBrowserImport compares summary markers to the local cache before the +// browser requests details. A fresh scan always begins at zero after restart; +// that is safe because this plan eliminates already-cached unchanged details. +func PlanBrowserImport(cacheRoot string, summaries []json.RawMessage) (BrowserImportPlan, error) { + return PlanBrowserImportWithForce(cacheRoot, summaries, false) +} + +// PlanBrowserImportWithForce requests every supplied conversation detail when +// force is true. It is used for explicit repair imports after users have +// permanently deleted previously imported Claude.ai sessions. +func PlanBrowserImportWithForce(cacheRoot string, summaries []json.RawMessage, force bool) (BrowserImportPlan, error) { + if err := os.MkdirAll(cacheRoot, 0o700); err != nil { + return BrowserImportPlan{}, fmt.Errorf("create Claude cache: %w", err) + } + manifest, err := readManifest(filepath.Join(cacheRoot, cacheManifestName)) + if err != nil { + return BrowserImportPlan{}, err + } + state, err := ReadSyncState(cacheRoot) + if err != nil { + return BrowserImportPlan{}, err + } + if state.Status != "running" { + state = SyncState{Status: "running", StartedAt: time.Now().UTC()} + } + // Always serialise an array so desktop clients can safely call includes + // when a page has no changed conversations. + plan := BrowserImportPlan{ChangedIDs: []string{}, State: state} + for _, summary := range summaries { + id, updatedAt, err := conversationMarker(summary) + if err != nil { + state.Failed++ + continue + } + state.Scanned++ + if !force && manifest[id] == updatedAt && fileExists(conversationCachePath(cacheRoot, id)) { + plan.Unchanged++ + state.Skipped++ + continue + } + plan.ChangedIDs = append(plan.ChangedIDs, id) + state.Changed++ + } + plan.State = state + if err := writeSyncState(cacheRoot, state); err != nil { + return BrowserImportPlan{}, err + } + return plan, nil +} + +// RecordBrowserImportProgress persists completed ingestion progress. The +// caller supplies importer counts after the database transaction succeeds. +func CompleteBrowserImport(cacheRoot string) (SyncState, error) { + return RecordBrowserImportProgress(cacheRoot, 0, 0, 0, true) +} + +func FailBrowserImport(cacheRoot string) (SyncState, error) { + state, err := ReadSyncState(cacheRoot) + if err != nil { + return SyncState{}, err + } + state.Status = "failed" + state.Failed++ + if err := writeSyncState(cacheRoot, state); err != nil { + return SyncState{}, err + } + return state, nil +} + +func RecordBrowserImportProgress(cacheRoot string, fetched, imported, failed int, finished bool) (SyncState, error) { + state, err := ReadSyncState(cacheRoot) + if err != nil { + return SyncState{}, err + } + state.Fetched += fetched + state.Imported += imported + state.Failed += failed + if finished { + state.Status = "completed" + } + if err := writeSyncState(cacheRoot, state); err != nil { + return SyncState{}, err + } + return state, nil +} + +// PreparedImport contains a Claude.ai export-compatible JSON payload. Call +// Commit only after that payload was successfully ingested; otherwise the next +// sync will refetch changed conversations rather than falsely marking them +// current. +type PreparedImport struct { + ExportJSON []byte + Downloaded int + Unchanged int + manifest map[string]string + manifestPath string +} + +// BrowserConversation is one authenticated browser fetch. The server receives +// this data only after Tauri fetched it from the isolated Claude webview; it +// never contains browser credentials. +type BrowserConversation struct { + Summary json.RawMessage `json:"summary"` + Conversation json.RawMessage `json:"conversation"` +} + +// PrepareBrowserImport updates the content cache from browser-fetched +// conversations and produces the same export-compatible payload as +// Browser imports deliberately have no HTTP or Keychain dependency. +func PrepareBrowserImport(cacheRoot string, conversations []BrowserConversation) (PreparedImport, error) { + return PrepareBrowserImportWithForce(cacheRoot, conversations, false) +} + +// PrepareBrowserImportWithForce saves and imports only the supplied batch. +// The durable cache is never replayed wholesale, avoiding quadratic imports. +func PrepareBrowserImportWithForce(cacheRoot string, conversations []BrowserConversation, force bool) (PreparedImport, error) { + if err := os.MkdirAll(cacheRoot, 0o700); err != nil { + return PreparedImport{}, fmt.Errorf("create Claude cache: %w", err) + } + manifestPath := filepath.Join(cacheRoot, cacheManifestName) + manifest, err := readManifest(manifestPath) + if err != nil { + return PreparedImport{}, err + } + result := PreparedImport{manifest: manifest, manifestPath: manifestPath} + batch := make([]json.RawMessage, 0, len(conversations)) + for _, item := range conversations { + id, updatedAt, err := conversationMarker(item.Summary) + if err != nil { + return PreparedImport{}, fmt.Errorf("invalid Claude conversation summary: %w", err) + } + path := conversationCachePath(cacheRoot, id) + if !force && manifest[id] == updatedAt && fileExists(path) { + result.Unchanged++ + continue + } + conversation, err := mergeSummary(item.Conversation, item.Summary) + if err != nil { + return PreparedImport{}, fmt.Errorf("normalise Claude conversation %s: %w", id, err) + } + if err := atomicWriteJSON(path, json.RawMessage(conversation)); err != nil { + return PreparedImport{}, fmt.Errorf("cache Claude conversation %s: %w", id, err) + } + manifest[id] = updatedAt + result.Downloaded++ + normalized, err := normalizeConversation(conversation) + if err != nil { + return PreparedImport{}, fmt.Errorf("normalise Claude import conversation %s: %w", id, err) + } + batch = append(batch, normalized) + } + payload, err := json.Marshal(batch) + if err != nil { + return PreparedImport{}, fmt.Errorf("encode Claude import batch: %w", err) + } + result.ExportJSON = payload + return result, nil +} + +// Commit records the remote update timestamps after a successful archive write. +func (p PreparedImport) Commit() error { + if p.manifestPath == "" { + return nil + } + return atomicWriteJSON(p.manifestPath, p.manifest) +} + +func readManifest(path string) (map[string]string, error) { + data, err := os.ReadFile(path) + if os.IsNotExist(err) { + return make(map[string]string), nil + } + if err != nil { + return nil, fmt.Errorf("read Claude cache index: %w", err) + } + manifest := make(map[string]string) + if err := json.Unmarshal(data, &manifest); err != nil { + return nil, fmt.Errorf("decode Claude cache index: %w", err) + } + return manifest, nil +} + +func conversationMarker(raw json.RawMessage) (string, string, error) { + var summary struct { + UUID string `json:"uuid"` + UpdatedAt string `json:"updated_at"` + } + if err := json.Unmarshal(raw, &summary); err != nil { + return "", "", err + } + if !safeConversationID(summary.UUID) { + return "", "", fmt.Errorf("invalid conversation UUID") + } + return summary.UUID, summary.UpdatedAt, nil +} + +func mergeSummary(detail, summary json.RawMessage) (json.RawMessage, error) { + var conversation map[string]json.RawMessage + if err := json.Unmarshal(detail, &conversation); err != nil { + return nil, err + } + var listed map[string]json.RawMessage + if err := json.Unmarshal(summary, &listed); err != nil { + return nil, err + } + for _, key := range []string{"uuid", "name", "created_at", "updated_at"} { + if len(conversation[key]) == 0 || string(conversation[key]) == `""` || string(conversation[key]) == "null" { + conversation[key] = listed[key] + } + } + return json.Marshal(conversation) +} + +// normalizeConversation preserves the fields consumed by the existing Claude +// export parser. Private API payloads include regenerated branches; retain the +// active leaf only so a transcript does not show competing replies twice. +func normalizeConversation(raw json.RawMessage) (json.RawMessage, error) { + var conversation struct { + UUID string `json:"uuid"` + Name string `json:"name"` + CreatedAt string `json:"created_at"` + UpdatedAt string `json:"updated_at"` + CurrentLeaf string `json:"current_leaf_message_uuid"` + Messages []json.RawMessage `json:"chat_messages"` + } + if err := json.Unmarshal(raw, &conversation); err != nil { + return nil, err + } + if !safeConversationID(conversation.UUID) { + return nil, fmt.Errorf("missing conversation UUID") + } + if conversation.CreatedAt == "" || conversation.UpdatedAt == "" { + return nil, fmt.Errorf("missing conversation timestamps") + } + messages, err := activeMessages(conversation.Messages, conversation.CurrentLeaf) + if err != nil { + return nil, err + } + return json.Marshal(struct { + UUID string `json:"uuid"` + Name string `json:"name"` + CreatedAt string `json:"created_at"` + UpdatedAt string `json:"updated_at"` + Messages []json.RawMessage `json:"chat_messages"` + }{conversation.UUID, conversation.Name, conversation.CreatedAt, conversation.UpdatedAt, messages}) +} + +func activeMessages(messages []json.RawMessage, leaf string) ([]json.RawMessage, error) { + if leaf == "" { + return messages, nil + } + type node struct { + UUID string `json:"uuid"` + Parent string `json:"parent_message_uuid"` + Raw json.RawMessage + } + byID := make(map[string]node, len(messages)) + for _, raw := range messages { + var item node + if err := json.Unmarshal(raw, &item); err != nil { + return nil, err + } + item.Raw = raw + if item.UUID != "" { + byID[item.UUID] = item + } + } + chain := make([]json.RawMessage, 0, len(messages)) + seen := make(map[string]struct{}) + for cursor := leaf; cursor != ""; { + item, ok := byID[cursor] + if !ok { + return messages, nil + } + if _, exists := seen[cursor]; exists { + return nil, fmt.Errorf("conversation message tree contains a cycle") + } + seen[cursor] = struct{}{} + chain = append(chain, item.Raw) + cursor = item.Parent + } + for left, right := 0, len(chain)-1; left < right; left, right = left+1, right-1 { + chain[left], chain[right] = chain[right], chain[left] + } + return chain, nil +} + +func conversationCachePath(root, id string) string { + return filepath.Join(root, "conversations", id+".json") +} + +func safeConversationID(id string) bool { + return id != "" && !strings.ContainsAny(id, `/\\`) && id != "." && id != ".." +} + +func fileExists(path string) bool { + info, err := os.Stat(path) + return err == nil && !info.IsDir() +} + +func atomicWriteJSON(path string, value any) error { + data, err := json.Marshal(value) + if err != nil { + return err + } + if err := os.MkdirAll(filepath.Dir(path), 0o700); err != nil { + return err + } + tmp, err := os.CreateTemp(filepath.Dir(path), ".tmp-*") + if err != nil { + return err + } + tmpName := tmp.Name() + defer os.Remove(tmpName) + if err := tmp.Chmod(0o600); err != nil { + _ = tmp.Close() + return err + } + if _, err := tmp.Write(append(data, '\n')); err != nil { + _ = tmp.Close() + return err + } + if err := tmp.Close(); err != nil { + return err + } + return os.Rename(tmpName, path) +} diff --git a/internal/cloudsync/claudeai/import_test.go b/internal/cloudsync/claudeai/import_test.go new file mode 100644 index 0000000000..977fa83caf --- /dev/null +++ b/internal/cloudsync/claudeai/import_test.go @@ -0,0 +1,56 @@ +package claudeai + +import ( + "encoding/json" + "path/filepath" + "testing" + + "github.com/stretchr/testify/assert" + "github.com/stretchr/testify/require" +) + +func TestPlanBrowserImportSkipsCachedSummaryAndPersistsState(t *testing.T) { + t.Parallel() + root := t.TempDir() + summary := []byte(`{"uuid":"conversation-1","updated_at":"2025-01-02T00:00:00Z"}`) + first, err := PlanBrowserImport(root, []json.RawMessage{summary}) + require.NoError(t, err) + assert.Equal(t, []string{"conversation-1"}, first.ChangedIDs) + assert.Equal(t, 1, first.State.Scanned) + require.NoError(t, atomicWriteJSON(conversationCachePath(root, "conversation-1"), json.RawMessage(`{"uuid":"conversation-1"}`))) + require.NoError(t, atomicWriteJSON(filepath.Join(root, cacheManifestName), map[string]string{"conversation-1": "2025-01-02T00:00:00Z"})) + second, err := PlanBrowserImport(root, []json.RawMessage{summary}) + require.NoError(t, err) + assert.Empty(t, second.ChangedIDs) + assert.Equal(t, 1, second.Unchanged) + repair, err := PlanBrowserImportWithForce(root, []json.RawMessage{summary}, true) + require.NoError(t, err) + assert.Equal(t, []string{"conversation-1"}, repair.ChangedIDs) + assert.Zero(t, repair.Unchanged) + state, err := CompleteBrowserImport(root) + require.NoError(t, err) + assert.Equal(t, "completed", state.Status) + state, err = FailBrowserImport(root) + require.NoError(t, err) + assert.Equal(t, "failed", state.Status) + assert.Equal(t, 1, state.Failed) +} + +func TestPrepareBrowserImportUsesOnlySuppliedBatchAndForceOverwritesCache(t *testing.T) { + t.Parallel() + root := t.TempDir() + summary := json.RawMessage(`{"uuid":"conversation-1","name":"one","created_at":"2025-01-01T00:00:00Z","updated_at":"2025-01-02T00:00:00Z"}`) + first, err := PrepareBrowserImport(root, []BrowserConversation{{Summary: summary, Conversation: json.RawMessage(`{"uuid":"conversation-1","chat_messages":[{"uuid":"m1","sender":"human","text":"old","created_at":"2025-01-01T00:00:00Z"}]}`)}}) + require.NoError(t, err) + require.NoError(t, first.Commit()) + second, err := PrepareBrowserImportWithForce(root, []BrowserConversation{{Summary: summary, Conversation: json.RawMessage(`{"uuid":"conversation-1","chat_messages":[{"uuid":"m1","sender":"human","text":"new","created_at":"2025-01-01T00:00:00Z"}]}`)}}, true) + require.NoError(t, err) + assert.JSONEq(t, `[{"uuid":"conversation-1","name":"one","created_at":"2025-01-01T00:00:00Z","updated_at":"2025-01-02T00:00:00Z","chat_messages":[{"uuid":"m1","sender":"human","text":"new","created_at":"2025-01-01T00:00:00Z"}]}]`, string(second.ExportJSON)) +} + +func TestNormalizeConversationUsesActiveBranch(t *testing.T) { + t.Parallel() + got, err := normalizeConversation([]byte(`{"uuid":"conversation-1","created_at":"2025-01-01T00:00:00Z","updated_at":"2025-01-02T00:00:00Z","current_leaf_message_uuid":"selected","chat_messages":[{"uuid":"root"},{"uuid":"discarded","parent_message_uuid":"root"},{"uuid":"selected","parent_message_uuid":"root"}]}`)) + require.NoError(t, err) + assert.JSONEq(t, `{"uuid":"conversation-1","name":"","created_at":"2025-01-01T00:00:00Z","updated_at":"2025-01-02T00:00:00Z","chat_messages":[{"uuid":"root"},{"uuid":"selected","parent_message_uuid":"root"}]}`, string(got)) +} diff --git a/internal/cloudsync/claudeai/service.go b/internal/cloudsync/claudeai/service.go new file mode 100644 index 0000000000..79e09bd976 --- /dev/null +++ b/internal/cloudsync/claudeai/service.go @@ -0,0 +1,429 @@ +package claudeai + +import ( + "bytes" + "context" + "encoding/json" + "errors" + "fmt" + "os" + "path/filepath" + "sync" + "time" + + "go.kenn.io/agentsview/internal/cloudsync/transport" + "go.kenn.io/agentsview/internal/db" + "go.kenn.io/agentsview/internal/importer" +) + +const ( + Provider = "claude-ai" + OperationListConversations = "list_conversations" + OperationGetConversation = "get_conversation" +) + +type ListParams struct { + Offset int `json:"offset"` + Limit int `json:"limit"` +} +type DetailParams struct { + ConversationID string `json:"conversation_id"` +} + +type SyncMode string + +const ( + SyncIncremental SyncMode = "incremental" + SyncRepair SyncMode = "repair" +) + +type JobStatus struct { + ID string `json:"id"` + Status string `json:"status"` + Mode SyncMode `json:"mode"` + Scanned int `json:"scanned"` + Changed int `json:"changed"` + Fetched int `json:"fetched"` + Imported int `json:"imported"` + Updated int `json:"updated"` + Skipped int `json:"skipped"` + Failed int `json:"failed"` + Error string `json:"error,omitempty"` +} + +type ScheduleConfig struct { + Enabled bool `json:"enabled"` + IntervalMinutes int `json:"interval_minutes"` + LastStartedAt time.Time `json:"last_started_at,omitempty"` + LastCompletedAt time.Time `json:"last_completed_at,omitempty"` + LastError string `json:"last_error,omitempty"` +} + +type Service struct { + broker *transport.Broker + store db.Store + cacheRoot, machine string + mu sync.Mutex + status JobStatus + cancel context.CancelFunc + schedule ScheduleConfig + archiveWrite func(func() error) error + stopSchedule chan struct{} + closed bool + schedulerWG sync.WaitGroup + jobsWG sync.WaitGroup +} + +func NewService(broker *transport.Broker, store db.Store, cacheRoot, machine string, archiveWrite ...func(func() error) error) *Service { + s := &Service{broker: broker, store: store, cacheRoot: cacheRoot, machine: machine, status: JobStatus{Status: "idle"}, schedule: ScheduleConfig{IntervalMinutes: 360}, stopSchedule: make(chan struct{})} + if len(archiveWrite) > 0 { + s.archiveWrite = archiveWrite[0] + } + if raw, err := os.ReadFile(filepath.Join(cacheRoot, "schedule.json")); err == nil { + _ = json.Unmarshal(raw, &s.schedule) + } + if s.schedule.IntervalMinutes < 15 { + s.schedule.IntervalMinutes = 360 + } + s.schedulerWG.Add(1) + go func() { + defer s.schedulerWG.Done() + s.runSchedule() + }() + return s +} + +func (s *Service) Schedule() ScheduleConfig { s.mu.Lock(); defer s.mu.Unlock(); return s.schedule } +func (s *Service) ConfigureSchedule(config ScheduleConfig) error { + if config.IntervalMinutes < 15 || config.IntervalMinutes > 24*60 { + return errors.New("Claude schedule interval must be between 15 minutes and 24 hours") + } + s.mu.Lock() + s.schedule.Enabled = config.Enabled + s.schedule.IntervalMinutes = config.IntervalMinutes + config = s.schedule + s.mu.Unlock() + return atomicWriteJSON(filepath.Join(s.cacheRoot, "schedule.json"), config) +} +func (s *Service) runSchedule() { + ticker := time.NewTicker(30 * time.Second) + defer ticker.Stop() + for { + select { + case <-s.stopSchedule: + return + case <-ticker.C: + } + s.mu.Lock() + cfg := s.schedule + running := s.status.Status == "running" || s.status.Status == "cancelling" + due := cfg.LastStartedAt.IsZero() || time.Since(cfg.LastStartedAt) >= time.Duration(cfg.IntervalMinutes)*time.Minute + s.mu.Unlock() + if !cfg.Enabled || running || !due { + continue + } + if _, err := s.Start(context.Background(), SyncIncremental); err == nil { + s.mu.Lock() + s.schedule.LastStartedAt = time.Now().UTC() + s.mu.Unlock() + _ = atomicWriteJSON(filepath.Join(s.cacheRoot, "schedule.json"), s.Schedule()) + } + } +} + +// Close stops the scheduler and joins all Claude work before archive shutdown. +func (s *Service) Close() { + s.mu.Lock() + if s.closed { + s.mu.Unlock() + return + } + s.closed = true + close(s.stopSchedule) + if s.cancel != nil && s.status.Status == "running" { + s.status.Status = "cancelling" + s.cancel() + } + s.mu.Unlock() + + s.schedulerWG.Wait() + s.jobsWG.Wait() +} +func (s *Service) Status() JobStatus { s.mu.Lock(); defer s.mu.Unlock(); return s.status } + +var ErrSyncAlreadyRunning = errors.New("Claude sync already running") + +func (s *Service) Start(ctx context.Context, mode SyncMode) (JobStatus, error) { + if s.store != nil && s.store.ReadOnly() { + return JobStatus{}, errors.New("cloud import not available in read-only mode") + } + if mode != SyncIncremental && mode != SyncRepair { + return JobStatus{}, fmt.Errorf("unsupported Claude sync mode %q", mode) + } + s.mu.Lock() + defer s.mu.Unlock() + if s.closed { + return JobStatus{}, errors.New("Claude sync service is closed") + } + if s.status.Status == "running" || s.status.Status == "cancelling" { + return s.status, ErrSyncAlreadyRunning + } + // HTTP request contexts end as soon as the start endpoint returns. The job + // owns its cancellation lifecycle instead of inheriting that short-lived + // context; callers still receive immediate start acknowledgement. + _ = ctx + jobCtx, cancel := context.WithCancel(context.Background()) + s.cancel = cancel + s.status = JobStatus{ID: fmt.Sprintf("claude-%d", time.Now().UnixNano()), Status: "running", Mode: mode} + s.jobsWG.Add(1) + go func() { + defer s.jobsWG.Done() + s.run(jobCtx) + }() + return s.status, nil +} +func (s *Service) Cancel() bool { + s.mu.Lock() + defer s.mu.Unlock() + if s.cancel == nil || s.status.Status != "running" { + return false + } + s.status.Status = "cancelling" + s.cancel() + return true +} +func (s *Service) update(fn func(*JobStatus)) { s.mu.Lock(); defer s.mu.Unlock(); fn(&s.status) } +func (s *Service) fail(err error) { + if !errors.Is(err, context.Canceled) { + _, _ = FailBrowserImport(s.cacheRoot) + s.persistScheduleResult(err.Error()) + } + s.update(func(st *JobStatus) { + if errors.Is(err, context.Canceled) { + st.Status = "cancelled" + } else { + st.Status = "failed" + st.Error = err.Error() + st.Failed++ + } + s.cancel = nil + }) +} +func (s *Service) finish() { + if _, err := CompleteBrowserImport(s.cacheRoot); err != nil { + s.fail(err) + return + } + s.update(func(st *JobStatus) { st.Status = "completed"; s.cancel = nil }) + s.persistScheduleResult("") +} + +func (s *Service) persistScheduleResult(lastError string) { + s.mu.Lock() + s.schedule.LastCompletedAt = time.Now().UTC() + s.schedule.LastError = lastError + config := s.schedule + s.mu.Unlock() + _ = atomicWriteJSON(filepath.Join(s.cacheRoot, "schedule.json"), config) +} + +func (s *Service) run(ctx context.Context) { + defer func() { + if r := recover(); r != nil { + s.fail(fmt.Errorf("Claude sync panic: %v", r)) + } + }() + seen := map[string]struct{}{} + for offset := 0; ; { + response, err := s.request(ctx, OperationListConversations, ListParams{Offset: offset, Limit: 50}) + if err != nil { + s.fail(err) + return + } + if response.Status == 401 || response.Status == 403 { + s.fail(fmt.Errorf("Claude authentication expired (HTTP %d)", response.Status)) + return + } + if response.Status < 200 || response.Status >= 300 { + s.fail(fmt.Errorf("Claude list returned HTTP %d", response.Status)) + return + } + items, hasMore, err := decodeBrowserPage(response.Body) + if err != nil { + s.fail(err) + return + } + if len(items) == 0 { + s.finish() + return + } + fresh := make([]json.RawMessage, 0, len(items)) + for _, item := range items { + id, _, e := conversationMarker(item) + if e != nil { + s.update(func(st *JobStatus) { st.Failed++ }) + continue + } + if _, ok := seen[id]; ok { + continue + } + seen[id] = struct{}{} + fresh = append(fresh, item) + } + s.update(func(st *JobStatus) { st.Scanned += len(fresh) }) + plan, err := PlanBrowserImportWithForce(s.cacheRoot, fresh, s.Status().Mode == SyncRepair) + if err != nil { + s.fail(err) + return + } + s.update(func(st *JobStatus) { st.Changed += len(plan.ChangedIDs); st.Skipped += plan.Unchanged }) + wanted := make(map[string]struct{}, len(plan.ChangedIDs)) + for _, id := range plan.ChangedIDs { + wanted[id] = struct{}{} + } + batch := make([]BrowserConversation, 0, len(wanted)) + for _, summary := range fresh { + id, _, _ := conversationMarker(summary) + if _, ok := wanted[id]; !ok { + continue + } + detail, err := s.request(ctx, OperationGetConversation, DetailParams{ConversationID: id}) + if err != nil { + s.fail(err) + return + } + if detail.Status == 404 { + s.update(func(st *JobStatus) { st.Skipped++ }) + continue + } + if detail.Status < 200 || detail.Status >= 300 { + s.fail(fmt.Errorf("Claude conversation %s returned HTTP %d", id, detail.Status)) + return + } + batch = append(batch, BrowserConversation{Summary: summary, Conversation: detail.Body}) + s.update(func(st *JobStatus) { st.Fetched++ }) + } + if len(batch) > 0 { + if err := s.importBatch(ctx, batch, s.Status().Mode == SyncRepair); err != nil { + s.fail(err) + return + } + } + offset += len(items) + if !hasMore { + s.finish() + return + } + } +} +func (s *Service) request(ctx context.Context, operation string, params any) (transport.Response, error) { + for attempt := 0; attempt < 5; attempt++ { + response, err := s.broker.Do(ctx, Provider, operation, params) + if err != nil { + return transport.Response{}, err + } + if response.Error != "" || response.Status == 429 || response.Status >= 500 { + if attempt == 4 { + if response.Error != "" { + return response, errors.New(response.Error) + } + return response, fmt.Errorf("Claude %s returned HTTP %d", operation, response.Status) + } + delay := time.Duration(1< 0 { + return fmt.Errorf("Claude import failed for %d conversations", stats.Errors) + } + if err := prepared.Commit(); err != nil { + return err + } + if _, err := RecordBrowserImportProgress(s.cacheRoot, prepared.Downloaded, stats.Imported+stats.Updated, 0, false); err != nil { + return err + } + s.update(func(st *JobStatus) { + st.Imported += stats.Imported + st.Updated += stats.Updated + st.Skipped += stats.Skipped + }) + return nil + }) +} + +func (s *Service) withArchiveWrite(work func() error) error { + if s.archiveWrite != nil { + return s.archiveWrite(work) + } + return work() +} +func decodeBrowserPage(raw json.RawMessage) ([]json.RawMessage, bool, error) { + var body struct { + Conversations []json.RawMessage `json:"conversations"` + Items []json.RawMessage `json:"items"` + Data []json.RawMessage `json:"data"` + Results []json.RawMessage `json:"results"` + HasMore *bool `json:"has_more"` + } + if err := json.Unmarshal(raw, &body); err != nil { + return nil, false, err + } + items := body.Conversations + if items == nil { + items = body.Items + } + if items == nil { + items = body.Data + } + if items == nil { + items = body.Results + } + if items == nil { + return nil, false, errors.New("Claude list had no conversations") + } + return items, body.HasMore == nil || *body.HasMore, nil +} diff --git a/internal/cloudsync/claudeai/service_test.go b/internal/cloudsync/claudeai/service_test.go new file mode 100644 index 0000000000..2d17085400 --- /dev/null +++ b/internal/cloudsync/claudeai/service_test.go @@ -0,0 +1,290 @@ +package claudeai + +import ( + "context" + "encoding/json" + "errors" + "path/filepath" + "sync" + "testing" + "time" + + "go.kenn.io/agentsview/internal/cloudsync/transport" + "go.kenn.io/agentsview/internal/db" + + "github.com/stretchr/testify/assert" + "github.com/stretchr/testify/require" +) + +func TestDecodeBrowserPageHonorsExplicitHasMore(t *testing.T) { + t.Parallel() + items, hasMore, err := decodeBrowserPage(json.RawMessage(`{"conversations":[{"uuid":"one"}],"has_more":true}`)) + require.NoError(t, err) + assert.Len(t, items, 1) + assert.True(t, hasMore, "a short page must continue when Claude explicitly sets has_more") + + _, hasMore, err = decodeBrowserPage(json.RawMessage(`{"conversations":[{"uuid":"one"}],"has_more":false}`)) + require.NoError(t, err) + assert.False(t, hasMore) +} + +func TestServiceStartReturnsRunningJobIdempotently(t *testing.T) { + t.Parallel() + service := NewService(transport.NewBroker(), nil, t.TempDir(), "local") + t.Cleanup(service.Close) + + started, err := service.Start(context.Background(), SyncIncremental) + require.NoError(t, err) + attached, err := service.Start(context.Background(), SyncIncremental) + require.ErrorIs(t, err, ErrSyncAlreadyRunning) + assert.Equal(t, started.ID, attached.ID) + assert.Equal(t, "running", attached.Status) +} + +func TestServiceRejectsReadOnlyStore(t *testing.T) { + writable := serviceTestDB(t) + path := writable.Path() + require.NoError(t, writable.Close()) + readonly, err := db.OpenReadOnly(path) + require.NoError(t, err) + t.Cleanup(func() { require.NoError(t, readonly.Close()) }) + service := NewService(transport.NewBroker(), readonly, t.TempDir(), "local") + t.Cleanup(service.Close) + _, err = service.Start(context.Background(), SyncIncremental) + require.EqualError(t, err, "cloud import not available in read-only mode") +} + +func TestServiceRunsBeyondStartRequestContext(t *testing.T) { + t.Parallel() + broker := transport.NewBroker() + service := NewService(broker, nil, t.TempDir(), "local") + t.Cleanup(service.Close) + requestContext, cancel := context.WithCancel(context.Background()) + _, err := service.Start(requestContext, SyncIncremental) + require.NoError(t, err) + cancel() + + claimContext, claimCancel := context.WithTimeout(context.Background(), time.Second) + defer claimCancel() + request, err := broker.Claim(claimContext) + require.NoError(t, err) + require.NoError(t, broker.Complete(transport.Response{ID: request.ID, Lease: request.Lease, Status: 200, Body: json.RawMessage(`{"conversations":[],"has_more":false}`)})) + require.Eventually(t, func() bool { return service.Status().Status == "completed" }, time.Second, 10*time.Millisecond) +} + +func TestSchedulePersistsAndCloseStopsActiveJob(t *testing.T) { + root := t.TempDir() + service := NewService(transport.NewBroker(), nil, root, "local") + t.Cleanup(service.Close) + require.NoError(t, service.ConfigureSchedule(ScheduleConfig{Enabled: true, IntervalMinutes: 60})) + reloaded := NewService(transport.NewBroker(), nil, root, "local") + t.Cleanup(reloaded.Close) + assert.True(t, reloaded.Schedule().Enabled) + assert.Equal(t, 60, reloaded.Schedule().IntervalMinutes) + _, err := service.Start(context.Background(), SyncIncremental) + require.NoError(t, err) + service.Close() + assert.Equal(t, "cancelled", service.Status().Status, "Close must join the cancelled sync job") +} + +func TestDecodeBrowserPageFallsBackWhenHasMoreIsAbsent(t *testing.T) { + t.Parallel() + _, hasMore, err := decodeBrowserPage(json.RawMessage(`{"conversations":[{"uuid":"one"}]}`)) + require.NoError(t, err) + assert.True(t, hasMore) +} + +func TestServiceImportsExplicitlyPaginatedConversationPages(t *testing.T) { + broker := transport.NewBroker() + database := serviceTestDB(t) + hadFTS := database.HasFTS() + service := NewService(broker, database, t.TempDir(), "local") + t.Cleanup(service.Close) + + var mu sync.Mutex + var offsets []int + stop := serveTransport(t, broker, func(request transport.Request) transport.Response { + switch request.Operation { + case OperationListConversations: + var params ListParams + require.NoError(t, json.Unmarshal(request.Params, ¶ms)) + mu.Lock() + offsets = append(offsets, params.Offset) + mu.Unlock() + if params.Offset == 0 { + return transportResponse(request, 200, `{"conversations":[`+string(testSummary("one"))+`],"has_more":true}`) + } + assert.Equal(t, hadFTS, database.HasFTS(), "FTS availability must be restored before the next remote page is fetched") + return transportResponse(request, 200, `{"conversations":[`+string(testSummary("two"))+`],"has_more":false}`) + case OperationGetConversation: + var params DetailParams + require.NoError(t, json.Unmarshal(request.Params, ¶ms)) + return transportResponse(request, 200, string(testDetail(params.ConversationID, "content "+params.ConversationID))) + default: + return transportResponse(request, 400, `{}`) + } + }) + t.Cleanup(stop) + + _, err := service.Start(context.Background(), SyncIncremental) + require.NoError(t, err) + require.Eventually(t, func() bool { return service.Status().Status == "completed" }, 2*time.Second, 10*time.Millisecond) + assert.Equal(t, []int{0, 1}, offsets) + assert.Equal(t, 2, service.Status().Imported) + for _, id := range []string{"one", "two"} { + session, err := database.GetSession(context.Background(), "claude-ai:"+id) + require.NoError(t, err) + require.NotNil(t, session) + } +} + +func TestServiceFetchesOnlyChangedConversations(t *testing.T) { + root := t.TempDir() + cached := BrowserConversation{Summary: testSummary("unchanged"), Conversation: testDetail("unchanged", "cached")} + prepared, err := PrepareBrowserImport(root, []BrowserConversation{cached}) + require.NoError(t, err) + require.NoError(t, prepared.Commit()) + + broker := transport.NewBroker() + service := NewService(broker, serviceTestDB(t), root, "local") + t.Cleanup(service.Close) + var detailIDs []string + var mu sync.Mutex + stop := serveTransport(t, broker, func(request transport.Request) transport.Response { + switch request.Operation { + case OperationListConversations: + return transportResponse(request, 200, `{"conversations":[`+string(testSummary("unchanged"))+`,`+string(testSummary("changed"))+`],"has_more":false}`) + case OperationGetConversation: + var params DetailParams + require.NoError(t, json.Unmarshal(request.Params, ¶ms)) + mu.Lock() + detailIDs = append(detailIDs, params.ConversationID) + mu.Unlock() + return transportResponse(request, 200, string(testDetail(params.ConversationID, "fresh"))) + default: + return transportResponse(request, 400, `{}`) + } + }) + t.Cleanup(stop) + + _, err = service.Start(context.Background(), SyncIncremental) + require.NoError(t, err) + require.Eventually(t, func() bool { return service.Status().Status == "completed" }, 2*time.Second, 10*time.Millisecond) + assert.Equal(t, []string{"changed"}, detailIDs) + assert.Equal(t, 1, service.Status().Changed) + assert.Equal(t, 1, service.Status().Skipped) +} + +func TestServiceFailedImportDoesNotCommitManifest(t *testing.T) { + broker := transport.NewBroker() + root := t.TempDir() + service := NewService(broker, serviceTestDB(t), root, "local", func(func() error) error { + return errors.New("archive write failed") + }) + t.Cleanup(service.Close) + stop := serveTransport(t, broker, func(request transport.Request) transport.Response { + switch request.Operation { + case OperationListConversations: + return transportResponse(request, 200, `{"conversations":[`+string(testSummary("uncommitted"))+`],"has_more":false}`) + case OperationGetConversation: + return transportResponse(request, 200, string(testDetail("uncommitted", "content"))) + default: + return transportResponse(request, 400, `{}`) + } + }) + t.Cleanup(stop) + + _, err := service.Start(context.Background(), SyncIncremental) + require.NoError(t, err) + require.Eventually(t, func() bool { return service.Status().Status == "failed" }, 2*time.Second, 10*time.Millisecond) + manifest, err := readManifest(filepath.Join(root, cacheManifestName)) + require.NoError(t, err) + assert.NotContains(t, manifest, "uncommitted") +} + +func TestServiceRepairRestoresOnlyFetchedClaudeSession(t *testing.T) { + broker := transport.NewBroker() + database := serviceTestDB(t) + service := NewService(broker, database, t.TempDir(), "local") + t.Cleanup(service.Close) + mode := SyncIncremental + stop := serveTransport(t, broker, func(request transport.Request) transport.Response { + switch request.Operation { + case OperationListConversations: + return transportResponse(request, 200, `{"conversations":[`+string(testSummary("repairable"))+`],"has_more":false}`) + case OperationGetConversation: + content := "initial" + if mode == SyncRepair { + content = "repaired" + } + return transportResponse(request, 200, string(testDetail("repairable", content))) + default: + return transportResponse(request, 400, `{}`) + } + }) + t.Cleanup(stop) + + _, err := service.Start(context.Background(), SyncIncremental) + require.NoError(t, err) + require.Eventually(t, func() bool { return service.Status().Status == "completed" }, 2*time.Second, 10*time.Millisecond) + require.NoError(t, database.DeleteSession("claude-ai:repairable")) + require.True(t, database.IsSessionExcluded("claude-ai:repairable")) + require.NoError(t, database.UpsertSession(db.Session{ + ID: "codex:unrelated", Project: "test", Machine: "local", Agent: "codex", + })) + require.NoError(t, database.DeleteSession("codex:unrelated")) + require.True(t, database.IsSessionExcluded("codex:unrelated")) + + mode = SyncRepair + _, err = service.Start(context.Background(), SyncRepair) + require.NoError(t, err) + require.Eventually(t, func() bool { return service.Status().Status == "completed" }, 2*time.Second, 10*time.Millisecond) + assert.False(t, database.IsSessionExcluded("claude-ai:repairable")) + assert.True(t, database.IsSessionExcluded("codex:unrelated")) + session, err := database.GetSession(context.Background(), "claude-ai:repairable") + require.NoError(t, err) + require.NotNil(t, session) +} + +func serviceTestDB(t *testing.T) *db.DB { + t.Helper() + database, err := db.Open(filepath.Join(t.TempDir(), "sessions.db")) + require.NoError(t, err) + t.Cleanup(func() { require.NoError(t, database.Close()) }) + return database +} + +func serveTransport(t *testing.T, broker *transport.Broker, handler func(transport.Request) transport.Response) func() { + t.Helper() + done := make(chan struct{}) + go func() { + for { + ctx, cancel := context.WithTimeout(context.Background(), 25*time.Millisecond) + request, err := broker.Claim(ctx) + cancel() + if err != nil { + select { + case <-done: + return + default: + continue + } + } + response := handler(request) + _ = broker.Complete(response) + } + }() + return func() { close(done) } +} + +func transportResponse(request transport.Request, status int, body string) transport.Response { + return transport.Response{ID: request.ID, Lease: request.Lease, Status: status, Body: json.RawMessage(body)} +} + +func testSummary(id string) json.RawMessage { + return json.RawMessage(`{"uuid":"` + id + `","name":"` + id + `","created_at":"2025-01-01T00:00:00Z","updated_at":"2025-01-02T00:00:00Z"}`) +} + +func testDetail(id, text string) json.RawMessage { + return json.RawMessage(`{"uuid":"` + id + `","chat_messages":[{"uuid":"message-` + id + `","sender":"human","text":"` + text + `","created_at":"2025-01-01T00:00:00Z"}]}`) +} diff --git a/internal/cloudsync/transport/transport.go b/internal/cloudsync/transport/transport.go new file mode 100644 index 0000000000..d70b35348f --- /dev/null +++ b/internal/cloudsync/transport/transport.go @@ -0,0 +1,165 @@ +// Package transport brokers credential-free provider requests between the +// archive server and an authenticated desktop adapter. Providers define typed +// operations; adapters never receive arbitrary remote URLs. +package transport + +import ( + "context" + "crypto/rand" + "encoding/hex" + "encoding/json" + "errors" + "fmt" + "sync" + "time" +) + +var ErrUnavailable = errors.New("authenticated transport unavailable") + +// The adapter allows a browser fetch up to 45 seconds, then still needs time +// to deliver its result over loopback. Keep the broker lease comfortably longer +// so a valid response cannot be reissued mid-delivery. +var leaseTimeout = 75 * time.Second + +type Request struct { + ID string `json:"id"` + Provider string `json:"provider"` + Operation string `json:"operation"` + Params json.RawMessage `json:"params"` + Lease string `json:"lease"` +} + +type Response struct { + ID string `json:"id"` + Lease string `json:"lease"` + Status int `json:"status"` + Body json.RawMessage `json:"body"` + RetryAfter string `json:"retry_after,omitempty"` + Error string `json:"error,omitempty"` +} + +type Broker struct { + mu sync.Mutex + queue chan *pending + pending map[string]*pending +} + +type pending struct { + request Request + result chan Response + timer *time.Timer +} + +func NewBroker() *Broker { + return &Broker{queue: make(chan *pending, 128), pending: make(map[string]*pending)} +} + +func (b *Broker) Do(ctx context.Context, provider, operation string, params any) (Response, error) { + data, err := json.Marshal(params) + if err != nil { + return Response{}, fmt.Errorf("encode transport params: %w", err) + } + p := &pending{request: Request{ID: randomID(), Provider: provider, Operation: operation, Params: data, Lease: randomID()}, result: make(chan Response, 1)} + b.mu.Lock() + b.pending[p.request.ID] = p + b.mu.Unlock() + defer func() { + b.mu.Lock() + delete(b.pending, p.request.ID) + if p.timer != nil { + p.timer.Stop() + } + b.mu.Unlock() + }() + select { + case b.queue <- p: + case <-ctx.Done(): + return Response{}, ctx.Err() + } + select { + case result := <-p.result: + return result, nil + case <-ctx.Done(): + return Response{}, ctx.Err() + } +} + +// Claim waits for a request. A desktop adapter calls this only after an +// authenticated webview is available. +func (b *Broker) Claim(ctx context.Context) (Request, error) { + select { + case p := <-b.queue: + b.mu.Lock() + if _, ok := b.pending[p.request.ID]; !ok { + b.mu.Unlock() + return Request{}, ErrUnavailable + } + p.request.Lease = randomID() + if p.timer != nil { + p.timer.Stop() + } + lease := p.request.Lease + p.timer = time.AfterFunc(leaseTimeout, func() { b.requeueExpired(p, lease) }) + request := p.request + b.mu.Unlock() + return request, nil + case <-ctx.Done(): + return Request{}, ctx.Err() + } +} + +func (b *Broker) Complete(response Response) error { + // Lookup, lease validation, and timer cancellation must be atomic. Otherwise + // an expiry callback can requeue the same request between validation and send. + b.mu.Lock() + p, ok := b.pending[response.ID] + if !ok { + b.mu.Unlock() + return errors.New("unknown or expired transport request") + } + if response.Lease == "" || response.Lease != p.request.Lease { + b.mu.Unlock() + return errors.New("invalid transport request lease") + } + if p.timer != nil { + p.timer.Stop() + p.timer = nil + } + b.mu.Unlock() + select { + case p.result <- response: + return nil + default: + return errors.New("transport request already completed") + } +} + +func (b *Broker) requeueExpired(p *pending, lease string) { + b.mu.Lock() + if current, ok := b.pending[p.request.ID]; !ok || current != p || p.request.Lease != lease { + b.mu.Unlock() + return + } + p.request.Lease = "" + p.timer = nil + b.mu.Unlock() + select { + case b.queue <- p: + default: + go func() { b.queue <- p }() + } +} + +func (b *Broker) Available() bool { + // Availability is intentionally inferred by a claimed request completing; + // a connected adapter need not expose credentials to prove its state. + return b != nil +} + +func randomID() string { + var raw [16]byte + if _, err := rand.Read(raw[:]); err != nil { + return fmt.Sprintf("fallback-%d", time.Now().UnixNano()) + } + return hex.EncodeToString(raw[:]) +} diff --git a/internal/cloudsync/transport/transport_test.go b/internal/cloudsync/transport/transport_test.go new file mode 100644 index 0000000000..796c4610c6 --- /dev/null +++ b/internal/cloudsync/transport/transport_test.go @@ -0,0 +1,56 @@ +package transport + +import ( + "context" + "testing" + "time" + + "github.com/stretchr/testify/assert" + "github.com/stretchr/testify/require" +) + +func TestBrokerDeliversSingleLeasedResponse(t *testing.T) { + broker := NewBroker() + result := make(chan Response, 1) + go func() { + response, err := broker.Do(context.Background(), "provider", "list", map[string]int{"offset": 0}) + require.NoError(t, err) + result <- response + }() + ctx, cancel := context.WithTimeout(context.Background(), time.Second) + defer cancel() + request, err := broker.Claim(ctx) + require.NoError(t, err) + require.NoError(t, broker.Complete(Response{ID: request.ID, Lease: request.Lease, Status: 200, Body: []byte(`[]`)})) + assert.Equal(t, 200, (<-result).Status) + assert.Error(t, broker.Complete(Response{ID: request.ID, Lease: request.Lease, Status: 200})) +} + +func TestBrokerRequeuesExpiredLease(t *testing.T) { + previous := leaseTimeout + leaseTimeout = 10 * time.Millisecond + t.Cleanup(func() { leaseTimeout = previous }) + broker := NewBroker() + result := make(chan error, 1) + go func() { _, err := broker.Do(context.Background(), "provider", "detail", struct{}{}); result <- err }() + ctx, cancel := context.WithTimeout(context.Background(), time.Second) + defer cancel() + first, err := broker.Claim(ctx) + require.NoError(t, err) + second, err := broker.Claim(ctx) + require.NoError(t, err) + assert.Equal(t, first.ID, second.ID) + assert.NotEqual(t, first.Lease, second.Lease) + require.NoError(t, broker.Complete(Response{ID: second.ID, Lease: second.Lease, Status: 200, Body: []byte(`[]`)})) + assert.NoError(t, <-result) +} + +func TestBrokerRejectsInvalidLease(t *testing.T) { + broker := NewBroker() + go func() { _, _ = broker.Do(context.Background(), "provider", "detail", struct{}{}) }() + ctx, cancel := context.WithTimeout(context.Background(), time.Second) + defer cancel() + request, err := broker.Claim(ctx) + require.NoError(t, err) + assert.Error(t, broker.Complete(Response{ID: request.ID, Lease: "wrong"})) +} diff --git a/internal/db/db_test.go b/internal/db/db_test.go index 6b38a9c44c..0eb635befc 100644 --- a/internal/db/db_test.go +++ b/internal/db/db_test.go @@ -2578,6 +2578,22 @@ func TestDeleteSessions(t *testing.T) { assert.Equal(t, 0, deleted, "deleted empty") } +func TestRestoreExcludedSessions(t *testing.T) { + d := testDB(t) + for _, id := range []string{"claude-ai:one", "claude-ai:two", "other:one"} { + insertSession(t, d, id, "p") + } + _, err := d.DeleteSessions([]string{"claude-ai:one", "claude-ai:two", "other:one"}) + require.NoError(t, err) + + restored, err := d.RestoreExcludedSessions([]string{"claude-ai:one", "claude-ai:two", "missing"}) + require.NoError(t, err) + assert.Equal(t, 2, restored) + assert.False(t, d.IsSessionExcluded("claude-ai:one")) + assert.False(t, d.IsSessionExcluded("claude-ai:two")) + assert.True(t, d.IsSessionExcluded("other:one")) +} + func TestDeleteSessionNonExistentNoGhostExclusion(t *testing.T) { d := testDB(t) diff --git a/internal/db/sessions.go b/internal/db/sessions.go index 40120cb205..c83a2ee684 100644 --- a/internal/db/sessions.go +++ b/internal/db/sessions.go @@ -1222,6 +1222,35 @@ func (db *DB) IsSessionExcluded(id string) bool { } // IsSessionTrashed returns true if the session ID exists in the trash. +// RestoreExcludedSessions removes permanent deletion markers for the supplied +// session IDs. It does not restore a trashed row; callers must re-import a +// trusted source after explicitly requesting this recovery. +func (db *DB) RestoreExcludedSessions(ids []string) (int, error) { + if err := db.requireWritable(); err != nil { + return 0, err + } + if len(ids) == 0 { + return 0, nil + } + placeholders := strings.TrimSuffix(strings.Repeat("?,", len(ids)), ",") + args := make([]any, len(ids)) + for i, id := range ids { + args[i] = id + } + db.mu.Lock() + defer db.mu.Unlock() + result, err := db.getWriter().Exec( + "DELETE FROM excluded_sessions WHERE id IN ("+placeholders+")", args...) + if err != nil { + return 0, fmt.Errorf("restore excluded sessions: %w", err) + } + count, err := result.RowsAffected() + if err != nil { + return 0, fmt.Errorf("count restored excluded sessions: %w", err) + } + return int(count), nil +} + func (db *DB) IsSessionTrashed(id string) bool { var n int _ = db.getReader().QueryRow( diff --git a/internal/db/store.go b/internal/db/store.go index e4f05bbe08..5f087b5367 100644 --- a/internal/db/store.go +++ b/internal/db/store.go @@ -141,6 +141,7 @@ type Store interface { DeleteSessionIfTrashed(id string) (int64, error) ListTrashedSessions(ctx context.Context) ([]Session, error) EmptyTrash() (int, error) + RestoreExcludedSessions(ids []string) (int, error) // Upload (local-only; PG returns ErrReadOnly). UpsertSession(s Session) error diff --git a/internal/duckdb/stubs.go b/internal/duckdb/stubs.go index 21cd2e0cf8..2eef43ca49 100644 --- a/internal/duckdb/stubs.go +++ b/internal/duckdb/stubs.go @@ -21,6 +21,7 @@ func (s *Store) SoftDeleteSessions(_ []string) (int, error) { return func (s *Store) RestoreSession(_ string) (int64, error) { return 0, db.ErrReadOnly } func (s *Store) DeleteSessionIfTrashed(_ string) (int64, error) { return 0, db.ErrReadOnly } func (s *Store) EmptyTrash() (int, error) { return 0, db.ErrReadOnly } +func (s *Store) RestoreExcludedSessions(_ []string) (int, error) { return 0, db.ErrReadOnly } func (s *Store) UpsertSession(_ db.Session) error { return db.ErrReadOnly } func (s *Store) ReplaceSessionMessages(_ string, _ []db.Message) error { return db.ErrReadOnly } func (s *Store) WriteSessionBatchAtomic( diff --git a/internal/importer/importer.go b/internal/importer/importer.go index 2e1397b269..419fe85ce4 100644 --- a/internal/importer/importer.go +++ b/internal/importer/importer.go @@ -54,23 +54,30 @@ type ftsSuspender interface { // (no message work happened), restore() is a no-op. This // avoids the expensive FTS rebuild when re-importing an // unchanged archive. -type lazyFTS struct { +// LazyFTS suspends FTS maintenance across one or more imports. Cloud syncs +// reuse one instance across pages so the index is rebuilt once per job. +type LazyFTS struct { sus ftsSuspender dropped bool onIndexing func() } -func newLazyFTS( +func NewLazyFTS( store db.Store, onIndexing func(), -) *lazyFTS { +) *LazyFTS { s, ok := store.(ftsSuspender) if !ok || !store.HasFTS() { return nil } - return &lazyFTS{sus: s, onIndexing: onIndexing} + return &LazyFTS{sus: s, onIndexing: onIndexing} } -func (f *lazyFTS) suspend() { +// newLazyFTS is retained for package-local callers and tests. +func newLazyFTS(store db.Store, onIndexing func()) *LazyFTS { + return NewLazyFTS(store, onIndexing) +} + +func (f *LazyFTS) suspend() { if f == nil || f.dropped { return } @@ -81,7 +88,7 @@ func (f *lazyFTS) suspend() { f.dropped = true } -func (f *lazyFTS) restore() error { +func (f *LazyFTS) Restore() error { if f == nil || !f.dropped { return nil } @@ -91,6 +98,7 @@ func (f *lazyFTS) restore() error { if err := f.sus.RebuildFTS(); err != nil { return fmt.Errorf("rebuilding FTS index: %w", err) } + f.dropped = false return nil } @@ -106,13 +114,48 @@ func ImportClaudeAI( cb *ImportCallbacks, machine ...string, ) (stats ImportStats, retErr error) { - fts := newLazyFTS(store, cb.indexing) + fts := NewLazyFTS(store, cb.indexing) defer func() { - if err := fts.restore(); err != nil { + if err := fts.Restore(); err != nil { retErr = errors.Join(retErr, err) } }() + return importClaudeAIWithFTS(ctx, store, r, cb, fts, machine...) +} + +// ImportClaudeAIWithFTS imports a batch while sharing an FTS suspension with +// its caller. The caller must call Restore after its complete import job. +func ImportClaudeAIWithFTS( + ctx context.Context, + store db.Store, + r io.Reader, + cb *ImportCallbacks, + fts *LazyFTS, + machine ...string, +) (stats ImportStats, retErr error) { + if fts == nil { + var onIndexing func() + if cb != nil { + onIndexing = cb.indexing + } + fts = NewLazyFTS(store, onIndexing) + defer func() { + if err := fts.Restore(); err != nil { + retErr = errors.Join(retErr, err) + } + }() + } + return importClaudeAIWithFTS(ctx, store, r, cb, fts, machine...) +} +func importClaudeAIWithFTS( + ctx context.Context, + store db.Store, + r io.Reader, + cb *ImportCallbacks, + fts *LazyFTS, + machine ...string, +) (stats ImportStats, retErr error) { provider, ok := parser.NewProvider( parser.AgentClaudeAI, parser.ProviderConfig{}, ) @@ -178,7 +221,7 @@ func upsertConversation( ctx context.Context, store db.Store, result parser.ParseResult, - fts *lazyFTS, + fts *LazyFTS, ) (importStatus, error) { s := result.Session @@ -292,9 +335,9 @@ func ImportChatGPT( cb *ImportCallbacks, machine ...string, ) (stats ImportStats, retErr error) { - fts := newLazyFTS(store, cb.indexing) + fts := NewLazyFTS(store, cb.indexing) defer func() { - if err := fts.restore(); err != nil { + if err := fts.Restore(); err != nil { retErr = errors.Join(retErr, err) } }() diff --git a/internal/parser/claude_ai.go b/internal/parser/claude_ai.go index 9fd6112420..cd346df8ba 100644 --- a/internal/parser/claude_ai.go +++ b/internal/parser/claude_ai.go @@ -1,6 +1,7 @@ package parser import ( + "bytes" "encoding/json" "fmt" "io" @@ -29,9 +30,23 @@ type claudeAIMessage struct { // Block types: text, thinking, tool_use, tool_result, // voice_note, token_budget. type claudeAIBlock struct { - Type string `json:"type"` - Text string `json:"text"` - Thinking string `json:"thinking"` + Type string `json:"type"` + Text string `json:"text"` + Thinking string `json:"thinking"` + Raw json.RawMessage `json:"-"` +} + +// UnmarshalJSON retains every provider block verbatim. Claude regularly adds +// block types; rendering an unknown block must not turn into data loss. +func (b *claudeAIBlock) UnmarshalJSON(data []byte) error { + type plain claudeAIBlock + var decoded plain + if err := json.Unmarshal(data, &decoded); err != nil { + return err + } + decoded.Raw = append(decoded.Raw[:0], data...) + *b = claudeAIBlock(decoded) + return nil } // ClaudeAIExportParser is implemented by the Claude.ai import-only provider to @@ -120,10 +135,12 @@ func assembleClaudeAIContent( } var contentParts []string + hasTextBlock := false for _, b := range m.Content { switch b.Type { case "text": if b.Text != "" { + hasTextBlock = true contentParts = append(contentParts, b.Text) } case "thinking": @@ -132,11 +149,21 @@ func assembleClaudeAIContent( contentParts = append(contentParts, "[Thinking]\n"+b.Thinking+"\n[/Thinking]") } - // tool_use, tool_result, voice_note, token_budget - // are metadata blocks — skip for display content. + default: + // Persist non-text blocks in message content as typed JSON fallback. + // This is deliberately verbose: data stays searchable/exportable until + // the UI grows a dedicated renderer for this Claude block type. + var pretty bytes.Buffer + if json.Indent(&pretty, b.Raw, "", " ") == nil { + contentParts = append(contentParts, "[Claude block: "+b.Type+"]\n```json\n"+pretty.String()+"\n```") + } } } + if !hasTextBlock && m.Text != "" { + contentParts = append([]string{m.Text}, contentParts...) + } + if len(contentParts) == 0 { if len(attachmentParts) == 0 { return m.Text, hasThinking diff --git a/internal/parser/claude_ai_test.go b/internal/parser/claude_ai_test.go index 8d3c827fe9..b49ce01c08 100644 --- a/internal/parser/claude_ai_test.go +++ b/internal/parser/claude_ai_test.go @@ -162,7 +162,10 @@ func TestParseClaudeAIExport_ContentBlocks(t *testing.T) { // Message with tool_use/tool_result blocks should use // text blocks, not the truncated top-level text. - assert.Equal(t, "First part.\n\nSecond part.", msgs[0].Content) + assert.Contains(t, msgs[0].Content, "First part.") + assert.Contains(t, msgs[0].Content, "Second part.") + assert.Contains(t, msgs[0].Content, "[Claude block: tool_use]") + assert.Contains(t, msgs[0].Content, "[Claude block: tool_result]") assert.False(t, msgs[0].HasThinking) // Message with thinking block. @@ -274,11 +277,9 @@ func TestParseClaudeAIExport_AttachmentFallbackPaths(t *testing.T) { "Top-level text survives.\n\nattachment with no filename", msgs[0].Content, ) - assert.Equal( - t, - "Fallback text survives too.\n\nattachment after unsupported block", - msgs[1].Content, - ) + assert.Contains(t, msgs[1].Content, "Fallback text survives too.") + assert.Contains(t, msgs[1].Content, "attachment after unsupported block") + assert.Contains(t, msgs[1].Content, "[Claude block: tool_use]") } func TestParseClaudeAIExport_IgnoredAttachments(t *testing.T) { diff --git a/internal/postgres/store.go b/internal/postgres/store.go index 978cabfb36..722077bd0e 100644 --- a/internal/postgres/store.go +++ b/internal/postgres/store.go @@ -558,6 +558,12 @@ func (s *Store) ListTrashedSessions( return scanPGSessionRows(rows) } +// RestoreExcludedSessions is local-archive recovery state and is unavailable +// when serving through PostgreSQL. +func (s *Store) RestoreExcludedSessions(_ []string) (int, error) { + return 0, db.ErrReadOnly +} + // EmptyTrash permanently deletes every trashed session. func (s *Store) EmptyTrash() (int, error) { ctx := context.Background() diff --git a/internal/server/huma_route_groups.go b/internal/server/huma_route_groups.go index 369684d9cf..92098db7f2 100644 --- a/internal/server/huma_route_groups.go +++ b/internal/server/huma_route_groups.go @@ -26,6 +26,7 @@ func (s *Server) registerTypedAPIRoutes() { s.registerStarredRoutes() s.registerPinRoutes() s.registerImportRoutes() + s.registerCloudRoutes() s.registerAssetRoutes() s.registerEmbeddingsRoutes() } diff --git a/internal/server/huma_routes_cloud.go b/internal/server/huma_routes_cloud.go new file mode 100644 index 0000000000..66ab54bd88 --- /dev/null +++ b/internal/server/huma_routes_cloud.go @@ -0,0 +1,97 @@ +package server + +import ( + "context" + "errors" + "net/http" + "time" + + "go.kenn.io/agentsview/internal/cloudsync/claudeai" + "go.kenn.io/agentsview/internal/cloudsync/transport" +) + +func (s *Server) registerCloudRoutes() { + group := newRouteGroup(s.api, "/api/v1/cloud", "Cloud sources") + post(s, group, "/claude-ai/sync", "Start Claude.ai sync", s.humaStartClaudeAICloud) + get(s, group, "/claude-ai/status", "Get Claude.ai sync status", s.humaClaudeAICloudStatus) + post(s, group, "/claude-ai/cancel", "Cancel Claude.ai sync", s.humaCancelClaudeAICloud) + get(s, group, "/claude-ai/schedule", "Get Claude.ai sync schedule", s.humaClaudeAICloudSchedule) + post(s, group, "/claude-ai/schedule", "Configure Claude.ai sync schedule", s.humaConfigureClaudeAICloudSchedule) + post(s, group, "/transport/claim", "Claim authenticated cloud request", s.humaClaimCloudTransport) + post(s, group, "/transport/result", "Complete authenticated cloud request", s.humaCompleteCloudTransport) +} + +type cloudSyncInput struct { + Body struct { + Mode claudeai.SyncMode `json:"mode"` + } `contentType:"application/json"` +} + +type cloudSyncResponse struct{ Body claudeai.JobStatus } + +type cloudScheduleInput struct { + Body claudeai.ScheduleConfig `contentType:"application/json"` +} + +type cloudScheduleResponse struct{ Body claudeai.ScheduleConfig } + +type cloudTransportClaimResponse struct{ Body transport.Request } + +type cloudTransportResultInput struct { + Body transport.Response `contentType:"application/json"` +} + +type cloudTransportResultResponse struct{ Body struct{} } + +func (s *Server) humaStartClaudeAICloud(ctx context.Context, in *cloudSyncInput) (*cloudSyncResponse, error) { + if s.db.ReadOnly() { + return nil, apiError(http.StatusNotImplemented, "cloud import not available in read-only mode") + } + status, err := s.claudeSync.Start(ctx, in.Body.Mode) + if errors.Is(err, claudeai.ErrSyncAlreadyRunning) { + // Starting an already-running local job is idempotent. This also lets a + // freshly loaded settings page attach to a scheduled sync. + return &cloudSyncResponse{Body: status}, nil + } + if err != nil { + return nil, apiError(http.StatusConflict, err.Error()) + } + return &cloudSyncResponse{Body: status}, nil +} + +func (s *Server) humaClaudeAICloudStatus(_ context.Context, _ *struct{}) (*cloudSyncResponse, error) { + return &cloudSyncResponse{Body: s.claudeSync.Status()}, nil +} + +func (s *Server) humaCancelClaudeAICloud(_ context.Context, _ *struct{}) (*cloudSyncResponse, error) { + s.claudeSync.Cancel() + return &cloudSyncResponse{Body: s.claudeSync.Status()}, nil +} + +func (s *Server) humaClaudeAICloudSchedule(_ context.Context, _ *struct{}) (*cloudScheduleResponse, error) { + return &cloudScheduleResponse{Body: s.claudeSync.Schedule()}, nil +} + +func (s *Server) humaConfigureClaudeAICloudSchedule(_ context.Context, in *cloudScheduleInput) (*cloudScheduleResponse, error) { + if err := s.claudeSync.ConfigureSchedule(in.Body); err != nil { + return nil, apiError(http.StatusBadRequest, err.Error()) + } + return &cloudScheduleResponse{Body: s.claudeSync.Schedule()}, nil +} + +func (s *Server) humaClaimCloudTransport(ctx context.Context, _ *struct{}) (*cloudTransportClaimResponse, error) { + claimCtx, cancel := context.WithTimeout(ctx, 20*time.Second) + defer cancel() + request, err := s.claudeTransport.Claim(claimCtx) + if err != nil { + return nil, apiError(http.StatusNoContent, "no authenticated cloud request is pending") + } + return &cloudTransportClaimResponse{Body: request}, nil +} + +func (s *Server) humaCompleteCloudTransport(_ context.Context, in *cloudTransportResultInput) (*cloudTransportResultResponse, error) { + if err := s.claudeTransport.Complete(in.Body); err != nil { + return nil, apiError(http.StatusBadRequest, err.Error()) + } + return &cloudTransportResultResponse{}, nil +} diff --git a/internal/server/server.go b/internal/server/server.go index 7b063254ad..71540adc97 100644 --- a/internal/server/server.go +++ b/internal/server/server.go @@ -10,6 +10,7 @@ import ( "net/http" httppprof "net/http/pprof" "net/url" + "path/filepath" "sort" "strconv" "strings" @@ -19,6 +20,8 @@ import ( "github.com/danielgtaylor/huma/v2" "github.com/danielgtaylor/huma/v2/adapters/humago" + "go.kenn.io/agentsview/internal/cloudsync/claudeai" + "go.kenn.io/agentsview/internal/cloudsync/transport" "go.kenn.io/agentsview/internal/config" "go.kenn.io/agentsview/internal/db" "go.kenn.io/agentsview/internal/insight" @@ -62,18 +65,20 @@ const ( // Server is the HTTP server that serves the SPA and REST API. type Server struct { - mu gosync.RWMutex - cfg config.Config - db db.Store - engine *sync.Engine - onDemandEngine *sync.Engine - sessions service.SessionService - broadcaster *Broadcaster - mux *http.ServeMux - api huma.API - httpSrv *http.Server - version VersionInfo - dataDir string + mu gosync.RWMutex + cfg config.Config + db db.Store + engine *sync.Engine + onDemandEngine *sync.Engine + sessions service.SessionService + broadcaster *Broadcaster + mux *http.ServeMux + api huma.API + httpSrv *http.Server + version VersionInfo + dataDir string + claudeTransport *transport.Broker + claudeSync *claudeai.Service httpRemoteCleanupRegistry *remotesync.CleanupRegistry @@ -184,6 +189,7 @@ func New( s := &Server{ cfg: cfg, db: database, + claudeTransport: transport.NewBroker(), engine: engine, sessions: sessions, mux: http.NewServeMux(), @@ -205,6 +211,9 @@ func New( spaFS: dist, spaHandler: http.FileServerFS(dist), } + s.claudeSync = claudeai.NewService(s.claudeTransport, database, + filepath.Join(cfg.DataDir, "cloud-cache", "claude-ai"), cfg.LocalMachineName, + s.serializeArchiveWrite) for _, opt := range opts { opt(s) } @@ -1118,6 +1127,7 @@ func (s *Server) Shutdown(ctx context.Context) error { if engine != nil { engine.Close() } + s.claudeSync.Close() return err }