diff --git a/Cargo.lock b/Cargo.lock index c5717a9..76267a3 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -2926,6 +2926,7 @@ dependencies = [ "mtp", "rand 0.8.6", "rand_core 0.6.4", + "reqwest", "sha2 0.10.9", "tokio", "uuid", diff --git a/config.json b/config.json new file mode 100644 index 0000000..03a370e --- /dev/null +++ b/config.json @@ -0,0 +1 @@ +{"keyring":"BMAGjI76myen0MXtUh7jwKfZBzJNeBofMRcVcDTJVSZy5IXAiR+f9GrbpQdIdYMyEgSRlbr6oaiqIBkhaSJZkI6YaIUFyIn4aRvPhG//ET/1NwLLyLw/1ripyKgAzAEFC3lGaLn0U7AS93vPuTxxMm5koKR7hzG40rTjEkcu1w6j1lpMupmFKqwpyj4WC2koVybhln8G0MjO947yXE0Tk47GsnWr0SENM7b79AlXNZvL9AMsGQbdgYqjIbMRZWthKY4P6EiZoWKFuYJrixeaVWSr+mCACgRDaMT+os+w+jw56oNCpFbBuswFR3iF8I3p9JdbiTl0ya1HwrngZJlUub9KxRgfCz5ZRhjzqKZ7Kgwhy7S4Nx2EFH4UhEaf+Y3QrGHitn7YCZXB5U0YGT+cqYJe8br77KiLW5rvjG2RaoaoA4AF8jiJ2LjWo1UrZDWL+hmiISuyWy3jZE7njI6oqXTZMb5oMSnB5r4KCyy/nA4DEMiTcLIcRgyUAXx1sDGX2Qc4yXwXqoMxnDzlCBFrRwn3EFbLYRNGpUM4y8Az4AROtTwTJq714RU9mc9v2V7jW7uCxYCMWWG2N3RIcppCJRbuwy1nWEPxNGEs8lp0wgDaSU1exz4WBldSCLD/mcImlCNiC174FVBgfLjM7Hp8oKnhCSziaTmAQ8mm5300Spw+A66VvCvGxheosy9BuTCpBKe9lJ/XdLrfnHxud1CoYgh9ymtGXDfkZ2H1c70XyC70+US/mFhfK6osMMDqOsd5eIRF43E/0oQfUcL8CBxEFcbDRyumSbEZpz42hCb9ccCzJZf6x6Bd9mzUeZYBlibrHJ+TKA6I+quB9CD1dljIQjTxQyLC91vbayoPuQ9iF2DKoEnicCNhhzLzZHRqE2oIE3MirDsCmMGno0+Cs7aqF8mvC1jBuzIQwlotOW7CYLjORmWLOrbbGMfV+XtOogEDWVYrxSe3dUOZ2x3fcQzGiVXjGjv2eTYQgHafDJUs6hbtDCqKUYXjzILF7MkpGx0EGYGhEQBYkMfQ7GQ8fI1ppYbpVUYh6m6zIWnWXII29aZOGldKAwO6uSjBwDI9d7lN1JgVR35Iu4Wj5zJj6ro9QxJ/mnSTwINSdj4SRo7JKIgWgXkeoS5VGM4Tq4umnIiQPFPXObrXBhQH8Lezsqyq83qeQFydy699zMz/WACflUua8Hhlgk6lMV0VfB3yBoJQx1WRxhGQm44CeBgt18u/Wm+ZTFJ+IR/+4yUg+qRhuc74k83LsCorFSSfymE4UDY7XJbAuZgCFVTs6TjzVCS4+YMdC4Row2p83J+CsYUAYoCnV67GNKEtwxH7XAJZFyKkc3b/kEPffDYJ9Ep7yrEy5AEqaWaT7B8f6U/6Nw4feykfinaruFPGplduEY9P6sLQ4Fvn+wDtA5EMwbkXQLwJSkGv3GOHBWbIWY+HmWnUW6lGPJd9O8US88SFIG0e8BdW0I9IaiAuvCVgpw8ImLJDswndSGV4sJ0A0Un4qIQrzCcszAR+GKArhMXcGo1SSndo9MxROyg1uJFwfJM9BQMZJjAE2+7LaIyaAGIf+/zEGWBoFhPYQqXu13laMW/npdMRUR40YydPhBRuMkXMeiJaTDFYCYCLmHASQghali1XJXH2TMJOh59iC2RiuJeTAZD7IFy/gaQtOLCxoHHz6L0kODXCdJB0xkr8CETSkgf7fKKlxSC8+n31IqRCq1tsgqssahhRBRFy+aQ4NzhVypvdoKiFqo3XmbRPhl3QxGHweaw9PMcq55au2SuttA9kZUcRo0C9PAppO3kLkTSjNamLWBj36Lmuo2KNlQjm9AS2CCwixxOSIWBj+SNCiWl3TF/0uhb6NEiKk3hSHDYBmbT+mK7KbIMqWM/Yemw9SyOdlI3KiEGRuj2BWzoy+ARSGxyCZ4HxGryAaRtWo8AwSaKD+Dz97LcEYUTC5Ja4qE/2yLzNtVJimTiFpkJLaojVEieWyZOdks8smBwTqz9I1aiE4GqX0no1/C4aYxtauMSWGRh9KV/1dLnEcgRExTvf6GIzEH+/kEi4Rh7AJLGuvKZT9kwEm7Nb9Sw4oUmHlZhyh5Q/gGChQZ3cDDFLdi4TBkePuAk10Dz2u575m5GvwULlN0NaJ3hx46KvyWbaVw84pIYuNg1Dq1wXWx0dOM+djL7LyXrQvJaD+spC+b0+N0o0MYcmvHgsG0VYm4IORMR7qAZg3HiBzF8Sdp1+Fz6jJG/ffDJTSmPtYSQDTIwctXrtpmw2+wOVW3G/hsf8RS7kNVSFxUnzrKTjmjovCFOaO8KEigz/CpPre6iUiVHX+w9ZIMBX8zScuLoObArc9Bn3Ykt2IcCoJaPW94V6pcgCUlzGi5lWo4cAx7pcSlL0pTRkk5AbcYxq+74kNKm4FFwB+rfPnE1R2QxyWTbi1psFgHJGchjypQBGeU1fUg19RCfz5HAL5YCCuL+Z4a+QiIGbdBrgSzBgUI6Hg1oG+Kp7IqJehG5/ekuMZ3MtNo9f1MLB603zUmcV+CE9gliIGFtVpzCwNlyRIbICwpT/1kqEjHFiK2jh4AgaWJ8ZqWg/mb6/qBz2BsqEqEVm+AL8BQsBOjtxh8uD3Aoat7mSApxNgrKKSh0Y5Emv1woxhUjueSFu5zRBWayaYl9oSkbMGZTPMGAYyT3vVDn2PKWEV7t8ND27PBktUM6/8ifisKI/K8ZTqDG7N6qioz3iY6sCZ2+Fq6jUBjQwqx34qYfRmHXMsX6198Fv0rrG2WZmEWc6QTJG60S0Ez90oAf72LathCal5MMRICSWEsx8BJpugX0ccaEgxZ6yJ5N2wa4FEzvDZXxL1LAJY5A05TvEB7fNQUSS9TnSc0+tu88ZhMKYMQjpMjnncoH+iL/VIRrlC4rxeRshoxHQsDD/23GzWz+n9buPiFwkQ1T+eCc/Q0XRIgrqPFmlIizSfAPhaQcNdcmfRkIrdosFNXxg95ciKAD14mD6ynjRZwYLSC+YKgCXQm7g1mOu1aGUBJkcl1nlxrnVQFudlpq3+RPCgS8g6JCCxFHrkgoBeqWZS6KfoxiOdZm6uj8LIzah6bFn2BFKC4xTEDQRGwxCHBFF5DKqd69ExqTFcLWD53/d0HFLmBfv1V7ESnzinKK2mQGpqss0VwdVsy16tLv1EqgGjI76myen0MXtUh7jwKfZBzJNeBofMRcVcDTJVSZy5IXAiR+f9GrbpQdIdYMyEgSRlbr6oaiqIBkhaSJZkI6YaIUFyIn4aRvPhG//ET/1NwLLyLw/1ripyKgAzAEFC3lGaLn0U7AS93vPuTxxMm5koKR7hzG40rTjEkcu1w6j1lpMupmFKqwpyj4WC2koVybhln8G0MjO947yXE0Tk47GsnWr0SENM7b79AlXNZvL9AMsGQbdgYqjIbMRZWthKY4P6EiZoWKFuYJrixeaVWSr+mCACgRDaMT+os+w+jw56oNCpFbBuswFR3iF8I3p9JdbiTl0ya1HwrngZJlUub9KxRgfCz5ZRhjzqKZ7Kgwhy7S4Nx2EFH4UhEaf+Y3QrGHitn7YCZXB5U0YGT+cqYJe8br77KiLW5rvjG2RaoaoA4AF8jiJ2LjWo1UrZDWL+hmiISuyWy3jZE7njI6oqXTZMb5oMSnB5r4KCyy/nA4DEMiTcLIcRgyUAXx1sDGX2Qc4yXwXqoMxnDzlCBFrRwn3EFbLYRNGpUM4y8Az4AROtTwTJq714RU9mc9v2V7jW7uCxYCMWWG2N3RIcppCJRbuwy1nWEPxNGEs8lp0wgDaSU1exz4WBldSCLD/mcImlCNiC174FVBgfLjM7Hp8oKnhCSziaTmAQ8mm5300Spw+A66VvCvGxheosy9BuTCpBKe9lJ/XdLrfnHxud1CoYgh9ymtGXDfkZ2H1c70XyC70+US/mFhfK6osMMDqOsd5eIRF43E/0oQfUcL8CBxEFcbDRyumSbEZpz42hCb9ccCzJZf6x6Bd9mzUeZYBlibrHJ+TKA6I+quB9CD1dljIQjTxQyLC91vbayoPuQ9iF2DKoEnicCNhhzLzZHRqE2oIE3MirDsCmMGno0+Cs7aqF8mvC1jBuzIQwlotOW7CYLjORmWLOrbbGMfV+XtOogEDWVYrxSe3dUOZ2x3fcQzGiVXjGjv2eTYQgHafDJUs6hbtDCqKUYXjzILF7MkpGx0EGYGhEQBYkMfQ7GQ8fI1ppYbpVUYh6m6zIWnWXII29aZOGldKAwO6uSjBwDI9d7lN1JgVR35Iu4Wj5zJj6ro9QxJ/mnSTwINSdj4SRo7JKIgWgXkeoS5VGM4Tq4umnIiQPFPXObrXBhQH8Lezsqyq83qeQFydy699zMz/WACflUua8Hhlgk6lMV0VfB3yBoJQx1WRxhGQm44CeBgt18u/Wm+ZTFJ+IR/+4yUg+qRhuc74k83LsCorFSSfymE4UDY7XJbAuZgCFVTs6TjzVCS4+YMdC4Row2p83J+CsYUAYoCnV67GNKEtwxH7XAJZFyKkc3b/kEPffDYJ9Ep7yrEy5AEqaWaT7B8f6U/6Nw4feykfinaruFPGplduEY9P6sLQ4Fvn+wDtA5EMwbkXQLwJSkGv3GOHBWbIWY+HmWnUW6lGPJd9O8US88SFIG0e8BdW0I9IaiAuvCVgpw8ImLJDswndSGV4sJ0A0Un4qIQrzCcszAR+GKArhMXcGo1SSndo9MxROyg1uJFwfJM9BQMZJjAE2+7LaIyaAGIf+/zEGWBoFu+riy8UUa4SXlhFlkzX0LxFF3GjDM8Ig7Z2e5dQPjhTF1E2/gR7C+el3nZRt3rb4xZF9OMFK1MOEdxBoJorubl0sqA89NULY0ePfMRef3eqe4PsNwcCSlUyLzqWwM10uQegMln/lXI/TJfBe4wuQ0IrPW3GBv+60A1pKSEmgyeT9RLSupHZGxbZHz0C5wocLbcpelkae/caApr7vUp8Kt3npl89ZM3xVj6HSb/8UgolC3yXfwG/W9ssCzmtQ4vC6HRhzM29MGBjZxMeLbYY+jrbI23L/iIgnKVeCenTh3wIVcoEKfkRvs8sAj1EjicL/YMJ0InKoiOrFC4CjMcclcuMw30Y8FmUcXxZjqltqbw0lhyuGHJTcUdesostQBapvrDkXTRK6hViPJehixApUEeTTLSUJDu9Nw2/LmOTEkYn0+sVDmkyjyYvaPmH/pvzuDIizEnYkRR/PKopmXhtnuis2UEh9U57qkNfJc7aj/OL245w062GpHpd2ynssXfbq87mjo4kfXwWACVAuxf7RWlbDGTfEF2wYeY+EjaouU137mA66uWPFhE3gyd9pvr2BsljYmcn3h7YStAMVWPe2m1Rtqlcbaez+O0g8fdG2elJgMt3asOhOzgJPIi19wRORZjoPPqqlZiwS1mu2m4EX91hJK1ZUT/d9UckuhaFVSyKAFo4edEKS+AkGsB6OBPuuvwxzwgYCv11eXcX5evB+b0mpEB+KPj5wkjWQcK51w5/Ondv+YPzgv1QN5SJznlz3kTn22P+qT/Zohz8Esy8EHZclpQ+lg+aF78ybvq3A9LKqiABSIX6ojLlQ0zODOm8UeTg18y8j8ofn5ivy9z3D1ihX7wwkSenDkDK8Fi+9wKTconAfFBoH4LINDHwiQuUUQ6BgmD6xpB4df7UdKowuODZO7nOQRy4mQVyrJMlvf8Hm54zw4vURSI4T5kFhbBizqcWEVftE+dki4Qtxxtd0rbUeXsxfykdgDrfkiGs7IoJL+cFQrY4rea1650/NLvQqk5/BILe0JybYPkzJLlFqHrL/mX1nG1RZJMfq8ChH1WECs9cZehKRFg7umjCEBXsr7+P713HqtVti33HYlQIAJEawbsM5Fj56/nxQiKVabUAykQ5krDq5OhLtn/z2IQAexjHmBTmkSLVWGl8P7JY4BfFQeAL/mW5TxVxf1MR7UE/GY1Re214zkrX+aa5tdWq0QKRa1xm5iepHYKgbAeU3D2jOwVzfhMWs3vDe1o4HjQrgC2+065TpcY57kizDEr1sX1mDvpaZBpTMBvyVdwpQsdSSUK3lruHkvYW5ATV4sdhPUua2ATk/gAVHC2WkFcxxiesCqkhSlib6YgolYDzVJu21dtR1mxMJed8FADYB7ICL7JgFBpGcit0Da0Khs1DAfCq1VT1cgHvXymzsjQPUKCFsixNKbk7fzByrU2IYOokjPAi4vA9jwYMsQ6fedoM2uyQT2Ro4nJ1DhGH5DaGdsi3nwKs6ceOadRSc43uYoVEyzpFcQyHAqBMYnz5e6tSTCqynK5cd5EwzJl/ShHhfRNH3nWoFWnYYenD4SWDqvKmyvLuxq4gISAw3onDRJsYde9M6gGvG6lWZtH3g/94qqH2g6RugRM5n/+5riMiN5AYY5LCMgkHeM4eOur5HKymjw1EVGKnPsX2m00bSkWSrTwUiiv3ogjmzM+F52CU7+BjYhdnaGdW+3zDCxfL4LIC/10MJD4Btt5vsdZwnXQdDizb/OmuVkIp47fPDR2Ni8LlWB3HPLcSQiE2B9RBKSeDoC5rr02SvS6be9oB9gYigmEp3QkhxfNdszW0Fw7rIWmzRRumRid6zqJRK39gH1f6XHFIJkTluZPqMdbHG6ZFSMxevrpNwUsjO/R8FYUYLw6ntj4VhU54T3s6Zh/kYLKi+qiVwPkv7RchRxB1gsnUgGBhA5nuU1AHtoABZorru7Wfrc9PaYty6mPrbyNKE56OCkhL/U5xqcSVKaZuLij5Qc60KzN0ixumiCPmKWaCSwK5rrjdKjD0OhTz5gmFmmaEisuvh94o1N6tEe2vYFHVTQ0biDq8/AbmahaSLU1TAcCjIznvl77nLBG1MJdMIH4kzyIw7uhSz/m/aET88RxIvc399PuTlXyOuysCfPBig1AMiRnpBWPhIlh9KxFj53P8QtfZgfBrEU/bHe27bkxFWkwGUxJSKeqHcrt8H7NzqaQbw7vdGXB6TO5saqz8qX1oJUzF0Fiuq0ce2nLeC5tpGWNswQ/WHEJmyJPGUAGCYyWE88qi2AKmJ/XBgZnugSCmtHjDQ6iuHt3+a9ZWUEl0ZzEf03PgOHCu/fjr7yAlk62+s9CW5L6YWx2ZQlrbvKWik8N76JS5lerB+TYIPqy4uYw+1UVzVC2yjmuqdnBU1Mm922FQ+JdZ4tmHQpHdswXiQSwvZkLN3MAle2yRQbIgmBcAUaI+rFBw4wUYM3UtAjjQoyKA2ATtRs7539E7KjjgNmpeHY/8ezDH+VqHscdOp1vncoMWI7AGgdai9LRnsGttSMddGMr2tgizkL6ucvIX5P8Oi57actaPRjeMRIsSyop9HQblrQhffJ1y0GlKwxI1gZQCqf955zn5P1uGCHGSh07qvP0e1qySPOb6rGkNhRV5xpPQ+Ik9YJBtx23j6kGwRI0DBFucoHJzRXsdgZZdn7MJ//jztsNzgVpx65QT8AAEeEX4qHHR8AVDqq49MrCQihEAIFXhH997tCM+mQuv4UGlq0pcTU3lxXZU/iFXBcKGNM57ACCjaD7MAxirp1DcoqFvmBAoogV577OS9UCCo+4cgL0MewAgXmP0NeyVgq5Dg60GcRXJ5P+4ulE/8wQX5boc1bhFEYY="} \ No newline at end of file diff --git a/iota-storage/src/users/user_manager.rs b/iota-storage/src/users/user_manager.rs index 248881a..bf1dc00 100644 --- a/iota-storage/src/users/user_manager.rs +++ b/iota-storage/src/users/user_manager.rs @@ -15,7 +15,8 @@ static UNIQUE: Lazy> = Lazy::new(|| Mutex::new(false)); pub async fn load_from_tu(username: &str) -> Result<(), ()> { let file_content = load_file("", &format!("{}.tu", username)); let segments = file_content.split("::").collect::>(); - let uuid = segments[0].parse::().unwrap_or(0); + let (uuid_str, _omega_host) = segments[0].split_once('@').unwrap_or((segments[0], "")); + let uuid = uuid_str.parse::().unwrap_or(0); let b64_private_key = segments[1]; let keyring = keyring_from_base64(b64_private_key).unwrap(); diff --git a/iota-storage/src/util/config_util.rs b/iota-storage/src/util/config_util.rs index 2dbaada..74e75e8 100644 --- a/iota-storage/src/util/config_util.rs +++ b/iota-storage/src/util/config_util.rs @@ -50,6 +50,14 @@ impl ConfigUtil { self.config["port"].as_u16().unwrap_or(1984) } + pub fn get_omikron_host(&self) -> Option { + self.config["omikron_host"].as_str().map(String::from) + } + + pub fn get_omikron_port(&self) -> Option { + self.config["omikron_port"].as_u16() + } + pub fn get_keyring(&self) -> Option { self.config["keyring"].as_str().map(String::from) } diff --git a/iota.mk b/iota.mk new file mode 100644 index 0000000..6108356 Binary files /dev/null and b/iota.mk differ diff --git a/omikron-connector/Cargo.toml b/omikron-connector/Cargo.toml index 6337e17..caaab62 100644 --- a/omikron-connector/Cargo.toml +++ b/omikron-connector/Cargo.toml @@ -16,6 +16,7 @@ mtp = { git = "https://git.methanium.net/Methanium/mtp.git", features = [ dashmap = "6.2.1" json = "*" +reqwest = "0.13.2" tokio = { version = "1.50.0", features = ["full"] } uuid = { version = "*", features = ["v4"] } base64 = "0.22.1" diff --git a/omikron-connector/src/lib.rs b/omikron-connector/src/lib.rs index 5d9124d..22029a6 100644 --- a/omikron-connector/src/lib.rs +++ b/omikron-connector/src/lib.rs @@ -1,3 +1,4 @@ +pub mod omega_discovery; pub mod omikron_connection; pub mod ping_pong_task; pub mod user_ops; diff --git a/omikron-connector/src/omega_discovery.rs b/omikron-connector/src/omega_discovery.rs new file mode 100644 index 0000000..7957cf0 --- /dev/null +++ b/omikron-connector/src/omega_discovery.rs @@ -0,0 +1,88 @@ +use std::env; +use std::time::Duration; + +use mtp::crypto::PublicKeyBundle; + +const OMEGA_API_BASE_DEFAULT: &str = "https://tensamin.net:9188"; +const REQUEST_TIMEOUT: Duration = Duration::from_secs(10); + +pub struct OmikronEndpoint { + pub id: i64, + pub host: String, + pub port: u16, + pub public_key: PublicKeyBundle, +} + +fn api_base() -> String { + env::var("OMEGA_API_URL").unwrap_or_else(|_| OMEGA_API_BASE_DEFAULT.to_string()) +} + +/* The Omega host as stored in `.tu` files: no `https://` scheme, but with port. */ +pub fn omega_host() -> String { + api_base() + .trim_start_matches("https://") + .trim_start_matches("http://") + .to_string() +} + +/* `GET /api/get/omikron` - random connected Omikron. Used on first-ever run; + * the only discovery endpoint with a liveness guarantee. */ +pub async fn discover_random() -> Result { + fetch(&format!("{}/api/get/omikron", api_base())).await +} + +/* `GET /api/get/omikron/{iota_id}` - this Iota's primary Omikron. + * No liveness guarantee (may 404 after restart or point at a stale Omikron); + * fall back to `discover_random`. */ +pub async fn discover_primary(iota_id: u64) -> Result { + fetch(&format!("{}/api/get/omikron/{}", api_base(), iota_id)).await +} + +async fn fetch(url: &str) -> Result { + let client = reqwest::Client::builder() + .timeout(REQUEST_TIMEOUT) + .build() + .map_err(|e| format!("Failed to build HTTP client: {}", e))?; + + let body = client + .get(url) + .send() + .await + .map_err(|e| format!("Request to {} failed: {}", url, e))? + .text() + .await + .map_err(|e| format!("Failed to read response body from {}: {}", url, e))?; + + let json = json::parse(&body).map_err(|e| format!("Invalid JSON from {}: {}", url, e))?; + + if json["status"].as_str() != Some("success") { + return Err(format!( + "Omega returned status {:?} for {}", + json["status"].as_str(), + url + )); + } + + let id = json["id"] + .as_i64() + .ok_or_else(|| format!("Missing/invalid \"id\" in response from {}", url))?; + let host = json["ip_address"] + .as_str() + .ok_or_else(|| format!("Missing/invalid \"ip_address\" in response from {}", url))? + .to_string(); + let port = json["port"] + .as_u16() + .ok_or_else(|| format!("Missing/invalid \"port\" in response from {}", url))?; + let public_key_b64 = json["public_key"] + .as_str() + .ok_or_else(|| format!("Missing/invalid \"public_key\" in response from {}", url))?; + let public_key = PublicKeyBundle::from_base64(public_key_b64) + .map_err(|e| format!("Failed to decode public key from {}: {}", url, e))?; + + Ok(OmikronEndpoint { + id, + host, + port, + public_key, + }) +} diff --git a/omikron-connector/src/omikron_connection.rs b/omikron-connector/src/omikron_connection.rs index ecaeffd..bf8a570 100755 --- a/omikron-connector/src/omikron_connection.rs +++ b/omikron-connector/src/omikron_connection.rs @@ -10,9 +10,11 @@ use iota_util::crypto_helper::{self, keyring_from_base64}; use iota_util::crypto_util::{self}; use iota_util::file_util::{get_children, has_file, load_file, save_file}; use json::JsonValue; -use mtp::client::{Policy, Receiver, SendMode, Sender}; +use mtp::client::{Client, ClientConfig, Policy, Receiver, SendMode, Sender}; use mtp::codec::{CommunicationType, CommunicationValue, DataType, DataValue}; +use mtp::crypto::{Keyring, PublicKeyBundle}; use std::collections::HashMap; +use std::env; use std::sync::{Arc, LazyLock}; use std::time::{Duration, Instant, SystemTime, UNIX_EPOCH}; use tokio::sync::{Mutex, RwLock, mpsc, watch}; @@ -20,6 +22,8 @@ use tokio::task::JoinHandle; use tokio::time::sleep; use uuid::Uuid; +use crate::omega_discovery; + fn typed_container(items: Vec<(DataType, DataValue)>) -> DataValue { use mtp::type_map::{DataTypeId, TypeMap}; let tm = TypeMap::latest(); @@ -47,8 +51,8 @@ async fn is_read_receipts_enabled() -> bool { // Configuration // ============================================================================ -const OMIKRON_HOST_DEFAULT: &str = "tensamin.net"; -const OMIKRON_PORT_DEFAULT: u16 = 959; +const IOTA_KEYRING_PATH: &str = "iota.mk"; +const OMIKRON_PUBLIC_KEY_PATH: &str = "omikron.mpkb"; const RECONNECT_DELAY: Duration = Duration::from_secs(5); const MAX_RECONNECT_DELAY: Duration = Duration::from_secs(300); const CONNECTION_TIMEOUT: Duration = Duration::from_secs(10); @@ -106,8 +110,6 @@ pub struct OmikronConnection { state: Arc>, sender: Arc>>>, connection_loop_handle: Arc>>>, - host: String, - port: u16, pub last_ping: Arc>, heartbeat_handle: Arc>>>, pub connection_id: Uuid, @@ -120,18 +122,12 @@ pub struct OmikronConnection { impl OmikronConnection { pub fn new() -> Self { - Self::with_host(OMIKRON_HOST_DEFAULT, OMIKRON_PORT_DEFAULT) - } - - pub fn with_host(host: &str, port: u16) -> Self { let (shutdown_tx, _) = watch::channel(false); OmikronConnection { state: Arc::new(RwLock::new(ConnectionState::Disconnected)), sender: Arc::new(RwLock::new(None)), connection_loop_handle: Arc::new(Mutex::new(None)), - host: host.to_string(), - port, last_ping: Arc::new(Mutex::new(-1)), heartbeat_handle: Arc::new(Mutex::new(None)), connection_id: Uuid::new_v4(), @@ -249,12 +245,21 @@ impl OmikronConnection { *self.state.write().await = ConnectionState::Connecting; log_t!("omikron_connecting"); - let addr_str = format!("https://{}:{}/ws/iota/", self.host, self.port); + let keyring = self.load_or_migrate_keyring().await; - let (sender, mut receiver) = mtp::client::client::connect( - &addr_str, - None, - Policy { + let existing_iota_id = match CONFIG.read().await.get_iota_id() { + 0 => None, + id => Some(id as u64), + }; + + let (host, port, omikron_public_key) = + self.resolve_omikron_endpoint(existing_iota_id).await?; + + let addr_str = format!("https://{}:{}/ws/iota/", host, port); + + let client_config = ClientConfig::new(&addr_str) + .with_description("iota") + .with_policy(Policy { send_mode: SendMode::SingleStreamPerMessage, max_message_size: 1_000_000_000, close_frame_len: u32::MAX, @@ -269,40 +274,51 @@ impl OmikronConnection { max_transient_recv_errors: 20, transient_recv_backoff: Duration::from_millis(100), receiver_queue_capacity: 1000, - }, + }); + + let connection = match Client::auth_connect_or_register( + client_config, + existing_iota_id, + &keyring, + &omikron_public_key, ) .await - .map_err(|e| format!("Connection failed: {}", e))?; + { + Ok(connection) => connection, + Err(mtp::common::CommunicationError::AuthenticationFailed(reason)) => { + let reason = format!( + "Authentication failed: {}. Your Iota keys may be invalid or the private key has changed on the server.", + reason + ); + *self.reconnect_on_close.write().await = false; + *self.auth_failure.write().await = Some(reason.clone()); + *self.state.write().await = ConnectionState::Disconnected; + return Err(reason); + } + Err(e) => return Err(format!("Connection failed: {}", e)), + }; log_t!("omikron_connection_success"); - let sender_arc = Arc::new(sender); + if existing_iota_id.is_none() { + let mut conf_write = CONFIG.write().await; + conf_write.change("iota_id", JsonValue::from(connection.client_id as i64)); + conf_write.update(); + drop(conf_write); + log!("Registered with Iota-ID: {}", connection.client_id); + } + + let sender_arc = Arc::new(connection.sender); *self.sender.write().await = Some(sender_arc.clone()); - *self.state.write().await = ConnectionState::Connected { identified: false }; + *self.state.write().await = ConnectionState::Connected { identified: true }; // Start read loop + let mut receiver = connection.receiver; let read_self = self.clone(); let read_handle = tokio::spawn(async move { read_self.read_loop(&mut receiver).await; }); - // Handle registration/identification - self.handle_authentication().await; - - // Wait for identification to complete - if !self.await_identification(Duration::from_secs(30)).await { - *self.reconnect_on_close.write().await = false; - let reason = "Authentication failed: server did not accept the challenge. Your Iota keys may be invalid or the private key has changed on the server." - .to_string(); - *self.auth_failure.write().await = Some(reason.clone()); - - if let Some(sender) = self.sender.write().await.take() { - sender.close(); - } - *self.state.write().await = ConnectionState::Disconnected; - return Err(reason); - } - log_t!("omikron_authenticated"); // Start heartbeat @@ -341,93 +357,144 @@ impl OmikronConnection { } // ------------------------------------------------------------------------- - // Authentication (Registration/Identification) + // Identity (own Keyring, migrated from the legacy base64-in-config format) // ------------------------------------------------------------------------- - async fn handle_authentication(&self) { - let conf = CONFIG.read().await; - let iota_id = conf.get_iota_id(); - let keyring_b64 = conf.get_keyring(); + /* + * `iota.mk` is now the source of truth for this Iota's identity. A + * pre-existing base64 keyring in config.json (from before the MTP auth + * migration) is imported once so already-registered Iotas keep their + * identity, and mirrored back into config.json for older code paths + * that still read it directly. + */ + async fn load_or_migrate_keyring(&self) -> Keyring { + if let Ok(kr) = mtp::files::load_keyring(IOTA_KEYRING_PATH) { + return kr; + } + + let legacy = CONFIG.read().await.get_keyring(); + let keyring = legacy + .and_then(|b64| keyring_from_base64(&b64)) + .unwrap_or_else(crypto_helper::generate_keyring); + + if let Err(e) = mtp::files::save_keyring(&keyring, IOTA_KEYRING_PATH) { + log!("Failed to persist {}: {}", IOTA_KEYRING_PATH, e); + } + + let b64 = crypto_helper::keyring_to_base64(&keyring); + let mut conf = CONFIG.write().await; + conf.change("keyring", JsonValue::from(b64)); + conf.update(); drop(conf); - if iota_id == 0 { - log_t!("iota_register_new"); + keyring + } - let pub_key_b64 = if let Some(kr) = keyring_b64 { - if let Some(keyring) = keyring_from_base64(&kr) { - let bundle = keyring.public_key_bundle(); - crypto_helper::public_key_bundle_to_base64(&bundle) - } else { - let keyring = crypto_helper::generate_keyring(); - let kb64 = crypto_helper::keyring_to_base64(&keyring); - let bundle = keyring.public_key_bundle(); - let pk_b64 = crypto_helper::public_key_bundle_to_base64(&bundle); - let mut conf_write = CONFIG.write().await; - conf_write.change("keyring", JsonValue::from(kb64)); - conf_write.update(); - drop(conf_write); - pk_b64 - } - } else { - let keyring = crypto_helper::generate_keyring(); - let kb64 = crypto_helper::keyring_to_base64(&keyring); - let bundle = keyring.public_key_bundle(); - let pk_b64 = crypto_helper::public_key_bundle_to_base64(&bundle); - let mut conf_write = CONFIG.write().await; - conf_write.change("keyring", JsonValue::from(kb64)); - conf_write.update(); - drop(conf_write); - pk_b64 - }; + // ------------------------------------------------------------------------- + // Omikron discovery (via Omega's HTTP API, replacing the static + // host/port/public-key-file model) + // ------------------------------------------------------------------------- - let register_msg = CommunicationValue::new(CommunicationType::RegisterIota) - .add_typed_default(DataType::PublicKey, DataValue::Str(pub_key_b64)); - - let msg_id = register_msg.get_id(); - - WAITING_TASKS.insert( - msg_id, - WaitingTask { - task: Box::new(|selfc, cv| { - if !cv.is_type(CommunicationType::Success) { - return false; - } - - let iota_value = cv.get_data(DataType::IotaId); - let iota_id = iota_value.as_number().unwrap_or(0); - - if iota_id != 0 { - tokio::spawn(async move { - let mut conf_write = CONFIG.write().await; - conf_write.change("iota_id", JsonValue::from(iota_id as i64)); - conf_write.update(); - drop(conf_write); - log!("Registered with Iota-ID: {}", iota_id); - - // Send identification after registration - let identify_msg = - CommunicationValue::new(CommunicationType::Identification) - .add_typed_default( - DataType::IotaId, - DataValue::SignedNumber(iota_id as i128), - ); - selfc.send_message(&identify_msg).await; - }); - } else { - log!("Iota registration failed."); - } - true - }), - inserted_at: Instant::now(), - }, - ); - - self.send_message(®ister_msg).await; - } else { - let identify_msg = CommunicationValue::new(CommunicationType::Identification) - .add_typed_default(DataType::IotaId, DataValue::SignedNumber(iota_id as i128)); - self.send_message(&identify_msg).await; + /* + * Discovery runs fresh on every `connect_once()` attempt rather than once + * at construction, since a fixed `OmikronConnection` may need to move to + * a different Omikron across reconnects (e.g. after the sticky/primary + * Omikron dies). `OMIKRON_HOST`/`OMIKRON_PORT` remain as a manual + * override for local dev/testing against a hand-run Omikron without a + * live Omega. + * + * The fetched Omikron public key is pinned to `omikron.mpkb` (trust on + * first use): if a cached key exists and a fresh discovery response + * disagrees with it, the mismatch is logged loudly and the cached key is + * kept rather than silently trusting whatever Omega's HTTP API returned + * this time - the same trust boundary the previous manual-file-drop + * model had, just automated for the common case. + */ + async fn resolve_omikron_endpoint( + &self, + existing_iota_id: Option, + ) -> Result<(String, u16, PublicKeyBundle), String> { + if let (Ok(host), Ok(port_str)) = (env::var("OMIKRON_HOST"), env::var("OMIKRON_PORT")) { + let port: u16 = port_str + .parse() + .map_err(|_| format!("Invalid OMIKRON_PORT: {}", port_str))?; + let public_key = mtp::files::load_public_key_bundle(OMIKRON_PUBLIC_KEY_PATH) + .map_err(|e| { + format!( + "Failed to load Omikron public key bundle from {}: {}. Obtain {} from the Omikron operator and place it in the working directory.", + OMIKRON_PUBLIC_KEY_PATH, e, OMIKRON_PUBLIC_KEY_PATH + ) + })?; + return Ok((host, port, public_key)); } + + let cached_key = mtp::files::load_public_key_bundle(OMIKRON_PUBLIC_KEY_PATH).ok(); + let cached_host_port = { + let conf = CONFIG.read().await; + match (conf.get_omikron_host(), conf.get_omikron_port()) { + (Some(host), Some(port)) => Some((host, port)), + _ => None, + } + }; + + let discovered = match existing_iota_id { + Some(id) => match omega_discovery::discover_primary(id).await { + Ok(endpoint) => Some(endpoint), + Err(e) => { + log!( + "Sticky Omikron discovery failed ({}), falling back to a random Omikron", + e + ); + omega_discovery::discover_random().await.ok() + } + }, + None => omega_discovery::discover_random().await.ok(), + }; + + let (host, port, public_key) = if let Some(endpoint) = discovered { + match &cached_key { + Some(cached) if cached.as_bytes() != endpoint.public_key.as_bytes() => { + log!( + "Fetched Omikron public key differs from the cached {} - keeping the \ + cached key. Delete {} manually if this is an expected key rotation.", + OMIKRON_PUBLIC_KEY_PATH, + OMIKRON_PUBLIC_KEY_PATH + ); + (endpoint.host, endpoint.port, cached.clone()) + } + Some(cached) => (endpoint.host, endpoint.port, cached.clone()), + None => { + if let Err(e) = mtp::files::save_public_key_bundle( + &endpoint.public_key, + OMIKRON_PUBLIC_KEY_PATH, + ) { + log!("Failed to cache Omikron public key: {}", e); + } + (endpoint.host, endpoint.port, endpoint.public_key) + } + } + } else if let (Some(cached), Some((host, port))) = (&cached_key, &cached_host_port) { + log!( + "Omega discovery unreachable, falling back to last-known Omikron {}:{}", + host, + port + ); + (host.clone(), *port, cached.clone()) + } else { + return Err( + "Omega discovery failed and no cached Omikron address/key is available" + .to_string(), + ); + }; + + { + let mut conf = CONFIG.write().await; + conf.change("omikron_host", JsonValue::from(host.clone())); + conf.change("omikron_port", JsonValue::from(port)); + conf.update(); + } + + Ok((host, port, public_key)) } // ------------------------------------------------------------------------- @@ -504,11 +571,6 @@ impl OmikronConnection { return; } - if cv.is_type(CommunicationType::Challenge) { - self.handle_challenge(&cv).await; - return; - } - if cv.is_type(CommunicationType::AppIdentification) { let sender_id = cv.get_sender(); let app_identifier = cv @@ -768,26 +830,6 @@ impl OmikronConnection { return; } - if cv.is_type(CommunicationType::IdentificationResponse) { - match cv.get_data(DataType::Accepted).as_bool() { - Some(true) => { - let mut state = self.state.write().await; - if let ConnectionState::Connected { identified: _ } = *state { - *state = ConnectionState::Connected { identified: true }; - } - } - Some(false) => { - *self.auth_failure.write().await = Some( - "Server rejected the challenge response — your Iota keys may be invalid." - .to_string(), - ); - log_t!("omikron_auth_rejected"); - } - None => {} - } - return; - } - // ************************************************ // // Direct messages // // ************************************************ // @@ -1694,65 +1736,6 @@ impl OmikronConnection { } } - async fn handle_challenge(&self, cv: &CommunicationValue) { - let conf = CONFIG.read().await; - let Some(kr_str) = conf.get_keyring() else { - drop(conf); - log_t!("omikron_challenge_decryption_failed"); - *self.auth_failure.write().await = Some( - "Challenge decryption failed: no keyring configured on this Iota.".to_string(), - ); - return; - }; - drop(conf); - - let Some(_omikron_pub_key_bundle) = cv.get_data(DataType::PublicKey).as_str() else { - log_t!("omikron_challenge_decryption_failed"); - return; - }; - let Some(encrypted_challenge) = cv.get_data(DataType::Challenge).as_str() else { - log_t!("omikron_challenge_decryption_failed"); - return; - }; - - let Some(keyring) = keyring_from_base64(&kr_str) else { - log_t!("omikron_challenge_decryption_failed"); - return; - }; - - let solved_challenge = crypto_util::decrypt_challenge(encrypted_challenge, &keyring).ok(); - - if let Some(solved) = solved_challenge { - let response = CommunicationValue::new(CommunicationType::ChallengeResponse) - .with_id(cv.get_id()) - .add_typed_default(DataType::Challenge, DataValue::Str(solved)); - - self.send_message(&response).await; - } else { - log_t!("omikron_challenge_decryption_failed"); - *self.auth_failure.write().await = Some( - "Challenge decryption failed — your Iota keyring may not match the registered keys on the server." - .to_string(), - ); - } - } - - async fn await_identification(&self, timeout: Duration) -> bool { - let start = Instant::now(); - loop { - if self.state.read().await.is_identified() { - return true; - } - if self.auth_failure.read().await.is_some() { - return false; - } - if start.elapsed() >= timeout { - return false; - } - sleep(Duration::from_millis(100)).await; - } - } - // ------------------------------------------------------------------------- // Public API // ------------------------------------------------------------------------- diff --git a/omikron-connector/src/user_ops.rs b/omikron-connector/src/user_ops.rs index 2d98042..d02ffbe 100644 --- a/omikron-connector/src/user_ops.rs +++ b/omikron-connector/src/user_ops.rs @@ -9,6 +9,7 @@ use mtp::codec::{CommunicationType, CommunicationValue, DataType, DataValue}; use rand_core::{OsRng, RngCore}; use std::time::Duration; +use crate::omega_discovery; use crate::omikron_connection::OMIKRON_CONNECTION; pub async fn create_user(username: &str) -> (Option, Option) { @@ -86,7 +87,7 @@ pub async fn create_user(username: &str) -> (Option, Option save_file( "", &format!("{}.tu", username), - &format!("{}::{}", user_id, keyring_b64), + &format!("{}@{}::{}", user_id, omega_discovery::omega_host(), keyring_b64), ); add_user(user_profile.clone());