diff --git a/Cargo.lock b/Cargo.lock index 212b238..475e451 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -3111,12 +3111,10 @@ dependencies = [ name = "monas-content" version = "0.1.0" dependencies = [ - "aes", "aes-gcm", "axum 0.8.8", "base64 0.22.1", "chrono", - "ctr", "dyn-clone", "hex", "hkdf", diff --git a/docs/design.md b/docs/design.md index 255b640..e297fed 100644 --- a/docs/design.md +++ b/docs/design.md @@ -157,7 +157,7 @@ presentation/ Axum HTTP API (port: 4002) | 機能 | 実装 | |------|------| -| コンテンツ暗号化 | AES-256-CTR(IVランダム生成) | +| コンテンツ暗号化 | AES-256-GCM(AEAD、12バイトランダムnonce。保存形式は `nonce \|\| ciphertext \|\| tag`) | | 鍵生成・管理 | CEK(Content Encryption Key)の生成・保存・削除 | | コンテンツアドレッシング | SHA-256によるCID生成 | | 鍵共有 | HPKE(RFC 9180、DH-KEM P-256)によるCEKのラップ | @@ -169,7 +169,7 @@ presentation/ Axum HTTP API (port: 4002) domain/ Content, ContentId, Share, Permission, KeyEnvelope application/ ContentService(CRUD + fetch + reencrypt) ShareService(grant, revoke, unwrap_cek) -infrastructure/ AES-256-CTR, HPKE, Sled, monas-filesync +infrastructure/ AES-256-GCM, HPKE, Sled, monas-filesync presentation/ Axum HTTP API (port: 4001) ``` @@ -344,6 +344,41 @@ Token失効は`min_valid_issued_at`による時刻ベースで管理される。 ネットワークはビザンチン耐性を前提として設計されている。悪意のあるノードが参加してもコンテンツの暗号化によって内容の漏洩は防がれる。XOR距離によるランダムなノード選択が一定の保護を提供する。 +#### relay先の信頼度 + +コンテンツを保持しないノードがリクエストを受けた場合、実際のmemberへrelayする。このときのrelay先候補には**由来の異なる2種類**があり、扱いを分ける必要がある。 + +| 由来 | 内容 | 扱い | +|---|---|---| +| ローカルの`ContentNetwork`レコード | 自ノードをmemberとして名指しした`ContentCreated` / `ContentNetworkManagerAdded`イベント由来。**発行元は認証済みだが、member集合そのものは発行元の主張** | ローカルレコード | +| DHTの近傍探索 | `sha256(content_id)`に近いというだけ。コンテンツとの関連は何も示されていない | 未証明 | + +ここで「ローカルレコード」は**owner署名によるattestationではない**。イベントにowner署名は無く、member集合は発行元の自己申告である。検証されているのは**発行元**の方で、Gossipsubを`MessageAuthenticity::Signed` + `ValidationMode::Strict`で運用しているため著者フィールドは必須かつ署名検証済みであり、これを次の2点に束縛している。 + +- `ContentCreated`は、名乗っている`creator_node_id`本人からの発行でなければ拒否する +- member集合を変える`ContentNetworkManagerAdded` / `ContentNetworkManagerRemoved`は、こちらが保持しているそのネットワークの既存memberからの発行でなければ拒否する +- `ContentDeleted`は上記に加えて、名乗っている`deleted_by_node_id`本人からの発行であることも確認する。ただしこのイベントは発行元を*自分で*名乗るので、その照合だけでは「認証済みなら誰でも通る」ことにしかならない。ローカルレコードを消せるのは既存memberだけである +- `ContentUpdated`は、名乗っている`updated_node_id`本人からの発行でなければ拒否する + +なお束縛に使うのはGossipsubの`Message::source`(**発行元**)であって`propagation_source`(直前の転送元)ではない。meshは多段転送するため、転送元で判定すると正規の多段配送を落としつつ偽装を通してしまう。 + +**候補が返した401/403は、出自によらず早期打ち切りの根拠にはしない。** 権威にすると、DHTキーの近くにPeer IDを置いた1台が403を返すだけであらゆるread/writeを止められてしまう(可用性への攻撃)。また正規のmemberであっても、policyの複製が終わっていない部分同期状態なら403を返し得るため、健全なレプリカへのfailoverを潰さないためにも継続が必要である。ただし答えとしては保持し、他の候補から何も得られなければそれを返す。 + +ローカルレコード由来のmemberについては、以前は「実policyに対する評価結果だから」として打ち切っていた。**これは撤回した。** レコード自体が最初の1通で植え付けられる(照合すべき既存membershipが無いため受理せざるを得ない)以上、その競争に勝った攻撃者は候補リストに載り、その403で正規callerのreadを恒久的に止められる — 未証明ピアについて防いでいるのと同じ攻撃が、ローカルレコード経路でも成立してしまう。早期打ち切りを戻せるのはowner署名付きmembership(#63)が入ってからである。継続のコストは「本当に拒否された場合に残り候補ぶんの往復が増える」ことに限られ、可用性側に倒すのが正しい方向である(callerはどのみち拒否され、それが少し遅くなるだけ)。 + +一方で、**credentialは未証明の候補へ転送してよい。** relayはcallerのtokenとリクエスト署名をそのまま転送するが、これは設計どおりであり、認可判断はmember側が実policyに対して行う。転送しなければmember側で認可できない。 + +これが安全なのは、認可が**Proof of Possession**だからである。tokenは自己完結型の鍵ID(公開鍵そのもの)か委譲JWTで、後者の`aud`(宛先)もまた自己完結型の鍵IDである。リクエスト署名は**その`aud`の鍵に対して**検証されるため、tokenと署名の両方を傍受した相手も、`aud`の秘密鍵を持たない以上、新しいリクエストを作れない。傍受した署名そのものも操作・リソース・body digest・timestampに束縛されており、mutationについてはさらに使い切りである。したがって未証明の候補へ渡っても、その相手ができるのは「同じreadを鮮度窓の内に再実行する」ことに限られる — readは冪等で、しかもその候補はrelay経由で既に暗号文を見ているため、新たに得られる情報はない。 + +候補リストを未証明だからといって切り詰めることはしない。切り詰めれば正当なmemberへのfailoverが減って可用性が落ちる一方、DHT距離順で先頭に来る相手は上限があろうと credential を受け取るため、機密性は改善しないからである。 + +残る課題は**member discoveryそのもの**である。発行元の認証によって「無関係なノードが勝手にレコードを植え付ける」ことは防げるが、次の2つの穴はプロトコル変更(owner署名付きmembership)でしか塞げず、未実装である。 + +- そのcontentについて**最初の**レコードは、照合すべき既存membershipが無いため受け入れざるを得ない +- 認証済みのmemberであれば、任意のmember集合を主張できる + +したがって「このノードが本当にこのcontentのmemberである」ことを暗号学的に確認する仕組みは依然として無く、ローカルレコード / 未証明の区別はそこへ至るまでの近似にとどまる。 + --- ## 11. CRSLとCRDT diff --git a/monas-content/Cargo.toml b/monas-content/Cargo.toml index a9a8668..506f6ec 100644 --- a/monas-content/Cargo.toml +++ b/monas-content/Cargo.toml @@ -11,8 +11,6 @@ path = "src/lib.rs" [dependencies] monas-filesync = { path = "../monas-filesync", optional = true } aes-gcm = "0.10.3" -aes = "0.8" -ctr = "0.9" rand_core = { version = "0.6.4", features = ["std"] } rand = "0.8.5" chrono = { version = "0.4.40", features = ["serde"] } diff --git a/monas-content/src/application_service/content_service/service.rs b/monas-content/src/application_service/content_service/service.rs index a664bbd..3396cd7 100644 --- a/monas-content/src/application_service/content_service/service.rs +++ b/monas-content/src/application_service/content_service/service.rs @@ -237,6 +237,29 @@ where }) } + /// ローカルに保存された「暗号化済み」バイト列を取得する。 + /// + /// State Node 整合性検証用途:State Node が保持するのは SDK が送信した + /// 暗号文なので、復号せずローカル暗号文とバイト比較することで + /// 「State Node が改ざんされていない同一の暗号文を保持しているか」を + /// 確認できる。 + pub fn fetch_encrypted(&self, content_id: ContentId) -> Result, FetchError> { + let content = self + .content_repository + .find_by_id(&content_id) + .map_err(FetchError::Repository)? + .ok_or(FetchError::NotFound)?; + + if content.is_deleted() { + return Err(FetchError::Deleted); + } + + content + .encrypted_content() + .cloned() + .ok_or(FetchError::NotFound) + } + /// 外部でアンラップされた CEK と暗号化済みコンテンツを用いて復号するユースケース。 /// /// - 共有フロー(Share)で KeyEnvelope から CEK を取り出した後の復号処理を想定。 @@ -1111,6 +1134,47 @@ mod tests { assert_eq!(stored.content_status(), &ContentStatus::Active); } + #[test] + fn fetch_encrypted_returns_stored_ciphertext() { + let (repo, _storage) = TestContentRepository::new(false); + let (key_store, _key_storage) = TestKeyStore::new(false, false); + let service = build_service(repo, TestKeyGenerator, TestEncryptor, key_store); + + let created = service + .create(CreateContentCommand { + name: "test".into(), + path: "path.txt".into(), + raw_content: b"hello".to_vec(), + provider: None, + }) + .expect("create should succeed"); + + let encrypted = service + .fetch_encrypted(created.content_id.clone()) + .expect("fetch_encrypted should succeed"); + // Contract: exactly the ciphertext produced at create time — the same + // bytes create() sent to the state node, so the verify-integrity + // comparison holds. (The test encryptor may be identity, so comparing + // against the plaintext would be meaningless here.) + assert_eq!(encrypted, created.encrypted_content); + + // Round-trip sanity: the plaintext fetch decrypts the same bytes. + let fetched = service + .fetch(created.content_id, None) + .expect("fetch should succeed"); + assert_eq!(fetched.raw_content, b"hello".to_vec()); + } + + #[test] + fn fetch_encrypted_not_found_for_unknown_content() { + let (repo, _) = TestContentRepository::new(false); + let (key_store, _) = TestKeyStore::new(false, false); + let service = build_service(repo, TestKeyGenerator, TestEncryptor, key_store); + + let result = service.fetch_encrypted(ContentId::new("missing".to_string())); + assert!(matches!(result, Err(FetchError::NotFound))); + } + #[test] fn create_validation_error_when_name_is_empty() { let (repo, _) = TestContentRepository::new(false); diff --git a/monas-content/src/infrastructure/encryption.rs b/monas-content/src/infrastructure/encryption.rs index 10e7591..6084014 100644 --- a/monas-content/src/infrastructure/encryption.rs +++ b/monas-content/src/infrastructure/encryption.rs @@ -3,13 +3,10 @@ use crate::domain::content::encryption::{ }; use crate::domain::content::ContentError; -use aes::Aes256; -use ctr::cipher::{KeyIvInit, StreamCipher}; -use ctr::Ctr128BE; +use aes_gcm::aead::Aead; +use aes_gcm::{Aes256Gcm, KeyInit, Nonce}; use rand_core::{OsRng, RngCore}; -type Aes256Ctr = Ctr128BE; - /// Implementation for generating a CEK suitable for AES-256. /// /// Produces a 32-byte random key using an OS-backed cryptographically secure RNG. @@ -24,18 +21,23 @@ impl ContentEncryptionKeyGenerator for OsRngContentEncryptionKeyGenerator { } } -/// Content encryption/decryption implementation using AES-256-CTR. +/// Content encryption/decryption implementation using AES-256-GCM (AEAD). /// -/// - Encryption: generates a 16-byte random IV and returns a byte sequence in the form `[iv || ciphertext]`. -/// - Decryption: splits the first 16 bytes as the IV and uses the remaining bytes as the ciphertext for AES-CTR. -/// - Provides confidentiality only; no integrity/authentication (no MAC or AEAD). -/// In the future this may be replaced with an AEAD scheme such as AES-GCM to add integrity protection. -pub struct Aes256CtrContentEncryption; - -const IV_LEN: usize = 16; +/// - Encryption: generates a 12-byte random nonce and returns a byte sequence +/// in the form `[nonce || ciphertext || tag]` (the 16-byte authentication +/// tag is appended by AES-GCM). +/// - Decryption: splits the first 12 bytes as the nonce and authenticates the +/// remaining bytes before returning the plaintext. Any tampering with the +/// ciphertext (including a forged payload substituted by an untrusted node) +/// fails authentication and returns a `DecryptionError` instead of silently +/// yielding corrupted plaintext. +pub struct Aes256GcmContentEncryption; + +const NONCE_LEN: usize = 12; +const TAG_LEN: usize = 16; const KEY_LEN: usize = 32; -impl ContentEncryption for Aes256CtrContentEncryption { +impl ContentEncryption for Aes256GcmContentEncryption { fn encrypt( &self, key: &ContentEncryptionKey, @@ -48,21 +50,24 @@ impl ContentEncryption for Aes256CtrContentEncryption { key.0.len() ))); } - let mut iv = [0u8; IV_LEN]; - let mut rng = OsRng; - rng.fill_bytes(&mut iv); - - let mut buffer = plaintext.to_vec(); - let mut cipher = Aes256Ctr::new_from_slices(key.0.as_slice(), &iv).map_err(|_| { + let cipher = Aes256Gcm::new_from_slice(key.0.as_slice()).map_err(|_| { ContentError::EncryptionError( - "Invalid key or IV length for AES-256-CTR (expected 32-byte key, 16-byte IV)" - .into(), + "Invalid key length for AES-256-GCM (expected 32-byte key)".into(), ) })?; - cipher.apply_keystream(&mut buffer); - let mut result = Vec::with_capacity(IV_LEN + buffer.len()); - result.extend_from_slice(&iv); - result.extend_from_slice(&buffer); + + let mut nonce_bytes = [0u8; NONCE_LEN]; + let mut rng = OsRng; + rng.fill_bytes(&mut nonce_bytes); + let nonce = Nonce::from_slice(&nonce_bytes); + + let ciphertext = cipher + .encrypt(nonce, plaintext) + .map_err(|_| ContentError::EncryptionError("AES-256-GCM encryption failed".into()))?; + + let mut result = Vec::with_capacity(NONCE_LEN + ciphertext.len()); + result.extend_from_slice(&nonce_bytes); + result.extend_from_slice(&ciphertext); Ok(result) } @@ -75,26 +80,26 @@ impl ContentEncryption for Aes256CtrContentEncryption { ))); } - if data.len() <= IV_LEN { + if data.len() < NONCE_LEN + TAG_LEN { return Err(ContentError::DecryptionError( - "Ciphertext is too short to contain IV and data (must be longer than IV only)" - .into(), + "Ciphertext is too short to contain nonce and authentication tag".into(), )); } - let (iv_bytes, ciphertext) = data.split_at(IV_LEN); - - let mut buffer = ciphertext.to_vec(); - - let mut cipher = Aes256Ctr::new_from_slices(key.0.as_slice(), iv_bytes).map_err(|_| { + let (nonce_bytes, ciphertext) = data.split_at(NONCE_LEN); + let cipher = Aes256Gcm::new_from_slice(key.0.as_slice()).map_err(|_| { ContentError::DecryptionError( - "Invalid key or IV length for AES-256-CTR (expected 32-byte key, 16-byte IV)" - .into(), + "Invalid key length for AES-256-GCM (expected 32-byte key)".into(), ) })?; - cipher.apply_keystream(&mut buffer); - Ok(buffer) + cipher + .decrypt(Nonce::from_slice(nonce_bytes), ciphertext) + .map_err(|_| { + ContentError::DecryptionError( + "AES-256-GCM authentication failed: ciphertext is corrupted or tampered".into(), + ) + }) } } @@ -105,7 +110,7 @@ mod tests { #[test] fn encrypt_then_decrypt_round_trip() { let key = ContentEncryptionKey(vec![42u8; 32]); // fixed key only for testing - let encryptor = Aes256CtrContentEncryption; + let encryptor = Aes256GcmContentEncryption; let plaintext = b"Monas content encryption test".to_vec(); let ciphertext = encryptor @@ -113,8 +118,7 @@ mod tests { .expect("encryption should succeed"); assert_ne!(ciphertext, plaintext); - assert!(ciphertext.len() > plaintext.len()); - assert!(ciphertext.len() >= IV_LEN); + assert_eq!(ciphertext.len(), NONCE_LEN + plaintext.len() + TAG_LEN); let decrypted = encryptor .decrypt(&key, &ciphertext) @@ -126,7 +130,7 @@ mod tests { #[test] fn encrypt_fails_with_invalid_key_length() { let key = ContentEncryptionKey(vec![1u8; 16]); - let encryptor = Aes256CtrContentEncryption; + let encryptor = Aes256GcmContentEncryption; let plaintext = b"test".to_vec(); let result = encryptor.encrypt(&key, &plaintext); @@ -136,9 +140,9 @@ mod tests { #[test] fn decrypt_fails_with_invalid_key_length() { let key = ContentEncryptionKey(vec![1u8; 16]); - let encryptor = Aes256CtrContentEncryption; + let encryptor = Aes256GcmContentEncryption; - let dummy_ciphertext = vec![0u8; IV_LEN + 4]; + let dummy_ciphertext = vec![0u8; NONCE_LEN + TAG_LEN + 4]; let result = encryptor.decrypt(&key, &dummy_ciphertext); assert!(matches!(result, Err(ContentError::DecryptionError(_)))); @@ -147,21 +151,22 @@ mod tests { #[test] fn decrypt_fails_when_data_too_short() { let key = ContentEncryptionKey(vec![2u8; 32]); - let encryptor = Aes256CtrContentEncryption; + let encryptor = Aes256GcmContentEncryption; - let too_short = vec![0u8; IV_LEN]; // exactly IV length (no payload data) + // nonce + tag だけの長さ未満は即エラー + let too_short = vec![0u8; NONCE_LEN + TAG_LEN - 1]; let result = encryptor.decrypt(&key, &too_short); assert!(matches!(result, Err(ContentError::DecryptionError(_)))); - let even_shorter = vec![0u8; IV_LEN - 1]; + let even_shorter = vec![0u8; NONCE_LEN]; let result2 = encryptor.decrypt(&key, &even_shorter); assert!(matches!(result2, Err(ContentError::DecryptionError(_)))); } #[test] - fn encrypt_produces_different_ciphertexts_due_to_random_iv() { + fn encrypt_produces_different_ciphertexts_due_to_random_nonce() { let key = ContentEncryptionKey(vec![99u8; 32]); - let encryptor = Aes256CtrContentEncryption; + let encryptor = Aes256GcmContentEncryption; let plaintext = b"same plaintext".to_vec(); let c1 = encryptor @@ -177,7 +182,7 @@ mod tests { #[test] fn encrypt_then_decrypt_round_trip_large_plaintext() { let key = ContentEncryptionKey(vec![7u8; 32]); - let encryptor = Aes256CtrContentEncryption; + let encryptor = Aes256GcmContentEncryption; let size = 1024 * 1024; let mut plaintext = Vec::with_capacity(size); for i in 0..size { @@ -188,8 +193,7 @@ mod tests { .encrypt(&key, &plaintext) .expect("encryption should succeed for large plaintext"); - assert!(ciphertext.len() > plaintext.len()); - assert!(ciphertext.len() >= IV_LEN); + assert_eq!(ciphertext.len(), NONCE_LEN + plaintext.len() + TAG_LEN); let decrypted = encryptor .decrypt(&key, &ciphertext) @@ -198,98 +202,88 @@ mod tests { assert_eq!(decrypted, plaintext); } - /// **Security vulnerability test**: This test passing demonstrates lack of integrity verification + /// **Integrity test**: AES-GCM detects any tampering of the ciphertext. /// - /// In AES-CTR mode, decryption succeeds even when ciphertext is tampered with. - /// Tampering goes undetected and incorrect data is returned. - /// - /// **Expected behavior (ideal)**: Tampered ciphertext should be detected and decryption should fail - /// **Current behavior (problem)**: Test passes = tampering undetected = security vulnerability + /// This is the counterpart of the old AES-CTR "vulnerability demo" test: + /// with CTR, a bit-flipped ciphertext decrypted "successfully" into + /// attacker-controlled plaintext. With GCM the authentication tag no + /// longer matches, so decryption MUST fail. This is what protects the SDK + /// from forged payloads returned by untrusted state nodes (PR #54 review). #[test] - fn tampered_ciphertext_decrypts_successfully_but_returns_wrong_data() { + fn tampered_ciphertext_fails_authentication() { let key = ContentEncryptionKey(vec![42u8; 32]); - let encryptor = Aes256CtrContentEncryption; + let encryptor = Aes256GcmContentEncryption; let plaintext = b"Secret message that should not be tampered with!"; - println!("\n========== TEST 1: Tampered Ciphertext =========="); - println!( - "Original Plaintext: {:?}", - String::from_utf8_lossy(plaintext) - ); - println!("Plaintext (hex): {}", hex::encode(plaintext)); - println!("Plaintext length: {} bytes", plaintext.len()); - - // Encrypt normally let mut ciphertext = encryptor .encrypt(&key, plaintext) .expect("encryption should succeed"); - println!("\n--- After Encryption ---"); - println!( - "Ciphertext length: {} bytes (IV: {} + data: {})", - ciphertext.len(), - IV_LEN, - ciphertext.len() - IV_LEN - ); - println!( - "IV (first 16 bytes): {}", - hex::encode(&ciphertext[0..IV_LEN]) - ); - println!( - "Encrypted data (hex): {}", - hex::encode(&ciphertext[IV_LEN..]) - ); + // Flip one bit in the encrypted payload (right after the nonce) + ciphertext[NONCE_LEN] ^= 0x11; - // Tamper with part of the ciphertext (modify the first byte after IV) - let original_byte = ciphertext[IV_LEN]; - println!("\n--- Before Tampering ---"); - println!("Byte at position [IV_LEN=16]: 0x{original_byte:02x} ({original_byte})"); - - ciphertext[IV_LEN] ^= 0x11; // Tampering by bit flipping - println!("\n--- After Tampering ---"); - let _tampered_byte = ciphertext[IV_LEN]; - println!( - "Byte at position [IV_LEN=16]: 0x{tampered:02x} ({tampered}) ← TAMPERED!", - tampered = ciphertext[IV_LEN] - ); - println!( - "Modified ciphertext (hex): {}", - hex::encode(&ciphertext[IV_LEN..]) + let result = encryptor.decrypt(&key, &ciphertext); + assert!( + matches!(result, Err(ContentError::DecryptionError(_))), + "tampered ciphertext must fail GCM authentication" ); - // Problem: Decryption "succeeds" even with tampered ciphertext - let decrypted = encryptor - .decrypt(&key, &ciphertext) - .expect("decryption 'succeeds' even with tampered data - THIS IS THE PROBLEM!"); - - println!("\n--- Decryption of Tampered Data ---"); - println!("WARNING: Decryption SUCCEEDED (this is the problem!)"); - println!("Decrypted text: {:?}", String::from_utf8_lossy(&decrypted)); - println!("Decrypted (hex): {}", hex::encode(&decrypted)); - - // Due to tampering, the decrypted result differs from the original first byte - println!("\n--- Verification ---"); - println!( - "Original 1st byte: 0x{:02x} ({})", - plaintext[0], plaintext[0] as char - ); - println!( - "Decrypted 1st byte: 0x{:02x} ({})", - decrypted[0], decrypted[0] as char - ); - assert_ne!(decrypted[0], plaintext[0]); - assert_ne!(&decrypted[..], &plaintext[..]); - println!("OK: Decrypted data DIFFERS from original (as expected from tampering)"); - - // Restoring the original byte allows correct decryption (proof that tampering was the cause) - println!("\n--- Restoring Original Ciphertext ---"); - ciphertext[IV_LEN] = original_byte; + // Restoring the original byte makes decryption succeed again, + // proving the failure above was caused by the tampering. + ciphertext[NONCE_LEN] ^= 0x11; let restored = encryptor .decrypt(&key, &ciphertext) - .expect("should decrypt"); - println!("Restored text: {:?}", String::from_utf8_lossy(&restored)); + .expect("untampered ciphertext should decrypt"); assert_eq!(&restored[..], &plaintext[..]); - println!("OK: After restoring byte, plaintext matches original"); - println!("========== END TEST 1 ==========\n"); + } + + /// Tampering with the nonce is also detected (the tag authenticates the + /// nonce implicitly via the keystream). + #[test] + fn tampered_nonce_fails_authentication() { + let key = ContentEncryptionKey(vec![42u8; 32]); + let encryptor = Aes256GcmContentEncryption; + let plaintext = b"nonce integrity"; + + let mut ciphertext = encryptor + .encrypt(&key, plaintext) + .expect("encryption should succeed"); + ciphertext[0] ^= 0x01; + + let result = encryptor.decrypt(&key, &ciphertext); + assert!(matches!(result, Err(ContentError::DecryptionError(_)))); + } + + /// Tampering with the authentication tag itself is detected. + #[test] + fn tampered_tag_fails_authentication() { + let key = ContentEncryptionKey(vec![42u8; 32]); + let encryptor = Aes256GcmContentEncryption; + let plaintext = b"tag integrity"; + + let mut ciphertext = encryptor + .encrypt(&key, plaintext) + .expect("encryption should succeed"); + let last = ciphertext.len() - 1; + ciphertext[last] ^= 0x80; + + let result = encryptor.decrypt(&key, &ciphertext); + assert!(matches!(result, Err(ContentError::DecryptionError(_)))); + } + + /// Decrypting with a wrong key fails instead of returning garbage. + #[test] + fn wrong_key_fails_authentication() { + let key = ContentEncryptionKey(vec![42u8; 32]); + let wrong_key = ContentEncryptionKey(vec![43u8; 32]); + let encryptor = Aes256GcmContentEncryption; + let plaintext = b"key binding"; + + let ciphertext = encryptor + .encrypt(&key, plaintext) + .expect("encryption should succeed"); + + let result = encryptor.decrypt(&wrong_key, &ciphertext); + assert!(matches!(result, Err(ContentError::DecryptionError(_)))); } } diff --git a/monas-content/src/presentation/mod.rs b/monas-content/src/presentation/mod.rs index ba8a5ca..adad938 100644 --- a/monas-content/src/presentation/mod.rs +++ b/monas-content/src/presentation/mod.rs @@ -12,7 +12,7 @@ use crate::{ application_service::{content_service::ContentService, share_service::ShareService}, infrastructure::{ content_id::Sha256ContentIdGenerator, - encryption::{Aes256CtrContentEncryption, OsRngContentEncryptionKeyGenerator}, + encryption::{Aes256GcmContentEncryption, OsRngContentEncryptionKeyGenerator}, key_store::InMemoryContentEncryptionKeyStore, key_wrapping::HpkeV1KeyWrapping, public_key_directory::InMemoryPublicKeyDirectory, @@ -36,7 +36,7 @@ struct AppState { Sha256ContentIdGenerator, MultiStorageRepository, OsRngContentEncryptionKeyGenerator, - Aes256CtrContentEncryption, + Aes256GcmContentEncryption, InMemoryContentEncryptionKeyStore, >, >, @@ -68,7 +68,7 @@ pub fn create_router() -> Router { content_id_generator: Sha256ContentIdGenerator, content_repository: content_repository.clone(), key_generator: OsRngContentEncryptionKeyGenerator, - encryptor: Aes256CtrContentEncryption, + encryptor: Aes256GcmContentEncryption, cek_store: cek_store.clone(), }; diff --git a/monas-sdk/src/controller/content.rs b/monas-sdk/src/controller/content.rs index e0e8058..5722d6e 100644 --- a/monas-sdk/src/controller/content.rs +++ b/monas-sdk/src/controller/content.rs @@ -24,7 +24,7 @@ use monas_content::domain::content::{Content, ContentEncryptionKey, StorageProvi use monas_content::domain::content_id::ContentId; use monas_content::infrastructure::{ content_id::Sha256ContentIdGenerator, - encryption::{Aes256CtrContentEncryption, OsRngContentEncryptionKeyGenerator}, + encryption::{Aes256GcmContentEncryption, OsRngContentEncryptionKeyGenerator}, MultiStorageRepository, }; @@ -38,7 +38,7 @@ pub(super) type ContentServiceInstance = ContentService< Sha256ContentIdGenerator, MultiStorageRepository, OsRngContentEncryptionKeyGenerator, - Aes256CtrContentEncryption, + Aes256GcmContentEncryption, DynCekStore, >; @@ -217,7 +217,7 @@ impl MonasController { .map(Some) } - fn prepare_state_node_metadata_auth( + pub(super) fn prepare_state_node_metadata_auth( &self, auth: Option<&StateNodeAuthContext>, operation: &str, diff --git a/monas-sdk/src/controller/mod.rs b/monas-sdk/src/controller/mod.rs index 4d02c27..4ca0023 100644 --- a/monas-sdk/src/controller/mod.rs +++ b/monas-sdk/src/controller/mod.rs @@ -262,14 +262,14 @@ impl MonasController { use monas_content::application_service::content_service::ContentService; use monas_content::infrastructure::{ content_id::Sha256ContentIdGenerator, - encryption::{Aes256CtrContentEncryption, OsRngContentEncryptionKeyGenerator}, + encryption::{Aes256GcmContentEncryption, OsRngContentEncryptionKeyGenerator}, }; ContentService { content_id_generator: Sha256ContentIdGenerator, content_repository, key_generator: OsRngContentEncryptionKeyGenerator, - encryptor: Aes256CtrContentEncryption, + encryptor: Aes256GcmContentEncryption, cek_store, } } diff --git a/monas-sdk/src/controller/share.rs b/monas-sdk/src/controller/share.rs index 7970d73..6e056dd 100644 --- a/monas-sdk/src/controller/share.rs +++ b/monas-sdk/src/controller/share.rs @@ -517,8 +517,14 @@ impl MonasController { } }; + // State Node は系列ID(remote_content_id)でコンテンツを管理する。 + // ローカル版IDしか送らないと State Node 側で未知のコンテンツ扱いになる。 + let state_node_content_id = input + .remote_content_id + .as_deref() + .unwrap_or(&input.content_id); if let Some(response) = self.send_update_to_state_node( - &input.content_id, + state_node_content_id, &reencryption.encrypted_content, auth, trace_id.clone(), diff --git a/monas-sdk/src/controller/state.rs b/monas-sdk/src/controller/state.rs index e1f8d80..91ea9a2 100644 --- a/monas-sdk/src/controller/state.rs +++ b/monas-sdk/src/controller/state.rs @@ -24,6 +24,29 @@ impl MonasController { None } + /// State Node の読み取り API 用の認証コンテキストを解決する。 + /// + /// 呼び出し元が Authorization を明示していればそのまま透過する。 + /// 無ければ書き込み系(create/update/delete)と同じく monas-account で + /// `read::` に署名し、`user:` トークンを組み立てる。 + /// State Node 側は読み取り時にこの署名メッセージを検証する + /// (`verify_read_access` → `verify_caller_signature("read", content_id, ..)`)。 + /// 署名を content_id にバインドすることで、relay 先ノード等に渡った署名を + /// 他コンテンツの読み取りに再利用されることを防ぐ。 + fn resolve_state_read_auth( + &self, + auth: Option<&StateNodeAuthContext>, + content_id: &str, + trace_id: &str, + ) -> Result, ApiResponse> { + match auth { + Some(ctx) if ctx.authorization.is_none() => { + self.prepare_state_node_metadata_auth(auth, "read", content_id, trace_id) + } + _ => Ok(auth.cloned()), + } + } + fn state_node_get_string( &self, url: &str, @@ -118,9 +141,17 @@ impl MonasController { return response; } + let auth = match self.resolve_state_read_auth::( + auth, + &input.content_id, + &trace_id, + ) { + Ok(resolved) => resolved, + Err(e) => return e, + }; let history = match self.get_state_node_history::( &input.content_id, - auth, + auth.as_ref(), trace_id.clone(), ) { Ok(h) => h, @@ -158,9 +189,17 @@ impl MonasController { return response; } + let auth = match self.resolve_state_read_auth::( + auth, + &input.content_id, + &trace_id, + ) { + Ok(resolved) => resolved, + Err(e) => return e, + }; let history = match self.get_state_node_history::( &input.content_id, - auth, + auth.as_ref(), trace_id.clone(), ) { Ok(h) => h, @@ -210,6 +249,16 @@ impl MonasController { ); } + let auth = match self.resolve_state_read_auth::( + auth, + &input.content_id, + &trace_id, + ) { + Ok(resolved) => resolved, + Err(e) => return e, + }; + let auth = auth.as_ref(); + let content_bytes = match URL_SAFE_NO_PAD.decode(&input.content) { Ok(b) => b, Err(e) => { @@ -264,13 +313,44 @@ impl MonasController { } }; - let valid = content_bytes == state_bytes; - let reason = if valid { - None + // State Node が保持するのは SDK が送信した「暗号文」なので、 + // local_content_id があればローカルに保存された暗号文とバイト比較する。 + // (平文 `content` と State Node のバイト列は一致し得ない。) + let (valid, reason) = if let Some(local_id) = input.local_content_id.as_deref() { + match self.content_service.fetch_encrypted( + monas_content::domain::content_id::ContentId::new(local_id.to_string()), + ) { + Ok(local_cipher) => { + if local_cipher == state_bytes { + (true, None) + } else { + ( + false, + Some(format!( + "state node ciphertext differs from local ciphertext (version={version_to_check})" + )), + ) + } + } + Err(e) => ( + false, + Some(format!( + "failed to load local ciphertext for {local_id}: {e}" + )), + ), + } } else { - Some(format!( - "content mismatch with state node (version={version_to_check})" - )) + let valid = content_bytes == state_bytes; + ( + valid, + if valid { + None + } else { + Some(format!( + "content mismatch with state node (version={version_to_check})" + )) + }, + ) }; ApiResponse::success( diff --git a/monas-sdk/src/models/share.rs b/monas-sdk/src/models/share.rs index 417f408..acb8935 100644 --- a/monas-sdk/src/models/share.rs +++ b/monas-sdk/src/models/share.rs @@ -75,7 +75,13 @@ pub struct DelegatedAccessToken { /// 共有取り消しリクエスト #[derive(Debug, Clone, Serialize, Deserialize)] pub struct RevokeShareInput { + /// SDK ローカルの版ID(ACL・CEK・再暗号化はローカルIDで処理される) pub content_id: String, + /// State Node へ送る系列ID。未指定の場合は `content_id` を使う(後方互換)。 + /// State Node はローカル版IDを知らないため、State Node に登録済みの + /// コンテンツでは必ず指定すること(`UpdateContentInput` と同じ区別)。 + #[serde(default, skip_serializing_if = "Option::is_none")] + pub remote_content_id: Option, /// 送信者の公開鍵(base64url) - sender_key_idを計算するために使用 pub sender_public_key: String, pub recipient_public_key: String, diff --git a/monas-sdk/src/models/state.rs b/monas-sdk/src/models/state.rs index 95aff8d..c50f0b4 100644 --- a/monas-sdk/src/models/state.rs +++ b/monas-sdk/src/models/state.rs @@ -54,6 +54,12 @@ pub struct VerifyIntegrityInput { pub content: String, #[serde(skip_serializing_if = "Option::is_none")] pub expected_version: Option, + /// SDK ローカルの版ID。指定すると、State Node が返す暗号文をローカルに + /// 保存された暗号文とバイト比較して検証する(State Node は暗号文を保持 + /// するため、平文である `content` とは直接比較できない)。未指定の場合は + /// 従来どおり `content` のバイト列と直接比較する。 + #[serde(default, skip_serializing_if = "Option::is_none")] + pub local_content_id: Option, } /// 整合性検証レスポンス diff --git a/monas-sdk/tests/share_controller_integration_test.rs b/monas-sdk/tests/share_controller_integration_test.rs index 91958e8..44f229a 100644 --- a/monas-sdk/tests/share_controller_integration_test.rs +++ b/monas-sdk/tests/share_controller_integration_test.rs @@ -186,6 +186,97 @@ async fn revoke_share_updates_state_node_version() { let revoke_response = controller.revoke_share( RevokeShareInput { content_id: created.content_id, + remote_content_id: None, + sender_public_key: sender.public_key, + recipient_public_key: recipient.public_key, + }, + None, + ); + assert!( + revoke_response.success, + "revoke_share should succeed: {:?}", + revoke_response.error + ); + update_mock.assert(); + + cleanup_content_artifacts(); +} + +#[tokio::test(flavor = "multi_thread")] +async fn revoke_share_syncs_state_node_by_remote_content_id() { + // The state node only knows the series id (remote_content_id), never the + // SDK-local version id. The post-revoke re-encryption PUT must therefore + // address the remote id when it is provided. + let _guard = acquire_test_lock(); + let mut server = Server::new_async().await; + let create_mock = server + .mock("POST", "/content") + .with_status(200) + .with_header("content-type", "application/json") + .with_body(r#"{"content_id":"remote-series-id"}"#) + .create_async() + .await; + let delegate_mock = server + .mock("POST", "/issuer/delegate") + .with_status(200) + .with_header("content-type", "application/json") + .with_body( + r#"{"delegated_token":"dummy.jwt.token","issued_at":1700000000,"expires_at":1700003600,"jti":"jti-3"}"#, + ) + .create_async() + .await; + // Only the remote id path is mocked: a PUT to any other path (e.g. the + // local content id) would fail the request and the assertion below. + let update_mock = server + .mock("PUT", "/content/remote-series-id") + .with_status(200) + .create_async() + .await; + + let controller = MonasController::with_urls(server.url(), server.url()); + + let sender = controller + .generate_keypair(GenerateKeypairInput { + key_type: KeyType::Secp256r1, + }) + .data + .expect("sender keypair should be generated"); + let recipient = controller + .generate_keypair(GenerateKeypairInput { + key_type: KeyType::Secp256r1, + }) + .data + .expect("recipient keypair should be generated"); + + let create_response = controller.create_content( + CreateContentInput { + content: URL_SAFE_NO_PAD.encode(b"revoke-remote-id-content"), + metadata: Some(ContentMetadata { + name: Some("revoke-remote.txt".to_string()), + content_type: Some("text/plain".to_string()), + created_at: None, + updated_at: None, + }), + }, + None, + ); + assert!(create_response.success, "create_content should succeed"); + let created = create_response.data.expect("create should return data"); + create_mock.assert(); + + let share_response = controller.share_content(ShareContentInput { + content_id: created.content_id.clone(), + sender_public_key: sender.public_key.clone(), + recipient_public_key: recipient.public_key.clone(), + permissions: vec![Permission::Write], + }); + assert!(share_response.success, "share_content should succeed"); + delegate_mock.assert(); + + let revoke_response = controller.revoke_share( + RevokeShareInput { + content_id: created.content_id, + remote_content_id: Some("remote-series-id".to_string()), sender_public_key: sender.public_key, recipient_public_key: recipient.public_key, }, @@ -280,6 +371,7 @@ async fn revoke_share_rolls_back_local_state_when_state_node_sync_fails() { let revoke_response = controller.revoke_share( RevokeShareInput { content_id: created.content_id.clone(), + remote_content_id: None, sender_public_key: sender.public_key.clone(), recipient_public_key: recipient.public_key.clone(), }, @@ -317,6 +409,7 @@ async fn revoke_share_rolls_back_local_state_when_state_node_sync_fails() { let second_revoke_response = controller.revoke_share( RevokeShareInput { content_id: created.content_id, + remote_content_id: None, sender_public_key: sender.public_key, recipient_public_key: recipient.public_key, }, @@ -414,6 +507,7 @@ async fn revoke_share_rollback_fires_on_inner_share_service_error() { let first = controller.revoke_share( RevokeShareInput { content_id: created.content_id.clone(), + remote_content_id: None, sender_public_key: sender.public_key.clone(), recipient_public_key: recipient.public_key.clone(), }, @@ -432,6 +526,7 @@ async fn revoke_share_rollback_fires_on_inner_share_service_error() { let second = controller.revoke_share( RevokeShareInput { content_id: created.content_id, + remote_content_id: None, sender_public_key: sender.public_key, recipient_public_key: recipient.public_key, }, diff --git a/monas-sdk/tests/state_controller_integration_test.rs b/monas-sdk/tests/state_controller_integration_test.rs index 6a4954f..706a339 100644 --- a/monas-sdk/tests/state_controller_integration_test.rs +++ b/monas-sdk/tests/state_controller_integration_test.rs @@ -121,6 +121,7 @@ async fn verify_integrity_rejects_stale_timestamp_with_unauthorized() { content_id: "test-content".into(), content: URL_SAFE_NO_PAD.encode(b"hello"), expected_version: Some("v1".into()), + local_content_id: None, }, Some(&auth), ); @@ -205,6 +206,7 @@ async fn verify_integrity_returns_api_error_when_history_cannot_be_fetched() { content_id: "test-content".into(), content: URL_SAFE_NO_PAD.encode(b"hello"), expected_version: None, + local_content_id: None, }, None, ); @@ -235,6 +237,7 @@ async fn verify_integrity_returns_api_error_when_version_cannot_be_fetched() { content_id: "test-content".into(), content: URL_SAFE_NO_PAD.encode(b"hello"), expected_version: Some("v1".into()), + local_content_id: None, }, None, ); @@ -265,6 +268,7 @@ async fn verify_integrity_keeps_false_only_for_actual_content_mismatch() { content_id: "test-content".into(), content: URL_SAFE_NO_PAD.encode(b"hello"), expected_version: Some("v1".into()), + local_content_id: None, }, None, ); @@ -303,6 +307,7 @@ async fn verify_integrity_returns_api_error_for_invalid_state_node_base64() { content_id: "test-content".into(), content: URL_SAFE_NO_PAD.encode(b"hello"), expected_version: Some("v1".into()), + local_content_id: None, }, None, ); diff --git a/monas-state-node/scripts/e2e-test.sh b/monas-state-node/scripts/e2e-test.sh index f9a7061..bead437 100755 --- a/monas-state-node/scripts/e2e-test.sh +++ b/monas-state-node/scripts/e2e-test.sh @@ -306,7 +306,7 @@ log_step "全ノードから content データを即座に取得できるか確 IMMEDIATE_MEMBERS=0 for port in 8080 8081 8082; do - generate_signature "$ACCOUNT1_PRIVATE_KEY" "read" "content" + generate_signature "$ACCOUNT1_PRIVATE_KEY" "read" "$CONTENT_ID" DATA_RESPONSE=$(curl -s "http://127.0.0.1:$port/content/$CONTENT_ID/data" \ -H "Authorization: Bearer $ACCOUNT1_KEY_ID" \ -H "X-Request-Signature: $LAST_SIGNATURE" \ @@ -445,7 +445,7 @@ echo "" updated_content_visible() { local p for p in 8080 8081 8082; do - generate_signature "$ACCOUNT1_PRIVATE_KEY" "read" "content" + generate_signature "$ACCOUNT1_PRIVATE_KEY" "read" "$CONTENT_ID" local resp resp=$(curl -s "http://127.0.0.1:$p/content/$CONTENT_ID/data" \ -H "Authorization: Bearer $ACCOUNT1_KEY_ID" \ @@ -464,7 +464,7 @@ poll_until 15 1 updated_content_visible || log_warn "15秒以内に更新の伝 log_step "各ノードでcontentデータを取得し、更新が反映されていることを確認" for port in 8080 8081 8082; do - generate_signature "$ACCOUNT1_PRIVATE_KEY" "read" "content" + generate_signature "$ACCOUNT1_PRIVATE_KEY" "read" "$CONTENT_ID" DATA_RESPONSE=$(curl -s "http://127.0.0.1:$port/content/$CONTENT_ID/data" \ -H "Authorization: Bearer $ACCOUNT1_KEY_ID" \ -H "X-Request-Signature: $LAST_SIGNATURE" \ @@ -483,7 +483,7 @@ done log_test "少なくとも1つのノードで更新データが取得できること" VERIFIED=false for port in 8080 8081 8082; do - generate_signature "$ACCOUNT1_PRIVATE_KEY" "read" "content" + generate_signature "$ACCOUNT1_PRIVATE_KEY" "read" "$CONTENT_ID" DATA_RESPONSE=$(curl -s "http://127.0.0.1:$port/content/$CONTENT_ID/data" \ -H "Authorization: Bearer $ACCOUNT1_KEY_ID" \ -H "X-Request-Signature: $LAST_SIGNATURE" \ diff --git a/monas-state-node/scripts/test-with-auth.sh b/monas-state-node/scripts/test-with-auth.sh index 370097a..30477a8 100755 --- a/monas-state-node/scripts/test-with-auth.sh +++ b/monas-state-node/scripts/test-with-auth.sh @@ -262,7 +262,7 @@ if [ -n "$CONTENT_ID" ]; then echo "" # CRDTデータの取得(認証ヘッダー付き) - generate_signature "$TEST_PRIVATE_KEY" "read" "content" + generate_signature "$TEST_PRIVATE_KEY" "read" "$CONTENT_ID" log_test "CRDTデータの取得" DATA_RESPONSE=$(curl -s -X GET "$BASE_URL/content/$CONTENT_ID/data" \ @@ -282,7 +282,7 @@ if [ -n "$CONTENT_ID" ]; then fi # CRDT履歴の取得(認証ヘッダー付き) - generate_signature "$TEST_PRIVATE_KEY" "read" "content" + generate_signature "$TEST_PRIVATE_KEY" "read" "$CONTENT_ID" log_test "CRDT履歴の取得" HIST_RESPONSE=$(curl -s -X GET "$BASE_URL/content/$CONTENT_ID/history" \ diff --git a/monas-state-node/src/application_service/content_sync_service.rs b/monas-state-node/src/application_service/content_sync_service.rs index 96644dd..a3c37d7 100644 --- a/monas-state-node/src/application_service/content_sync_service.rs +++ b/monas-state-node/src/application_service/content_sync_service.rs @@ -101,13 +101,25 @@ where } }; - // 2. Get local version to request only newer operations - let local_version = self + // 2. Get local version to request only newer operations. + // Only trust the local history if we actually hold the genesis node: + // crsl-lib's `linear_history` returns `[genesis]` even for content we + // hold nothing of (phantom history). Passing that as `since` makes + // providers skip the Create operation and we would never converge. + let has_genesis = self .crdt_repo - .get_history(genesis_cid) + .has_genesis(genesis_cid) .await - .ok() - .and_then(|h| h.last().cloned()); + .unwrap_or(false); + let local_version = if has_genesis { + self.crdt_repo + .get_history(genesis_cid) + .await + .ok() + .and_then(|h| h.last().cloned()) + } else { + None + }; // 3. Fetch operations from each member node for node_id in network.member_nodes() { @@ -365,6 +377,42 @@ mod tests { ) } + #[tokio::test] + async fn test_sync_requests_full_history_when_genesis_missing() { + // Phantom-history guard: when we do not actually hold the genesis + // node, the sync must request the FULL operation history + // (since=None). Passing the phantom [genesis] as `since` makes + // providers skip the Create operation and the node never converges. + let service = create_service_with_members("node-1", "content-1", vec!["node-2"], vec![]); + // MockContentRepository is empty -> has_genesis("content-1") == false. + + let _ = service.sync_from_peers("content-1").await.unwrap(); + + let since = service.peer_network.fetch_operations_since.lock().await; + assert_eq!(since.as_slice(), &[None]); + } + + #[tokio::test] + async fn test_sync_requests_incremental_when_genesis_present() { + let service = create_service_with_members("node-1", "content-1", vec!["node-2"], vec![]); + // Materialize the content locally so has_genesis == true. + service + .crdt_repo + .contents + .lock() + .await + .insert("content-1".to_string(), b"data".to_vec()); + service.crdt_repo.history.lock().await.insert( + "content-1".to_string(), + vec!["v1".to_string(), "v2".to_string()], + ); + + let _ = service.sync_from_peers("content-1").await.unwrap(); + + let since = service.peer_network.fetch_operations_since.lock().await; + assert_eq!(since.as_slice(), &[Some("v2".to_string())]); + } + #[tokio::test] async fn test_sync_from_peers_no_network() { let service = create_test_service("node-1"); diff --git a/monas-state-node/src/application_service/node.rs b/monas-state-node/src/application_service/node.rs index ee90b61..f83d07a 100644 --- a/monas-state-node/src/application_service/node.rs +++ b/monas-state-node/src/application_service/node.rs @@ -335,7 +335,9 @@ impl StateNode { let service_for_relay = self.service.clone(); let token_relay = token.clone(); tokio::spawn(async move { - use crate::infrastructure::network::libp2p_network::RelayRequestKind; + use crate::infrastructure::network::libp2p_network::{ + RelayOutcome, RelayRequestKind, + }; use crate::port::auth_token::AuthToken; tracing::info!("Started relay request handler"); loop { @@ -364,7 +366,7 @@ impl StateNode { timestamp, ) .await - .map(|_| ()) + .map(|_| RelayOutcome::Done) } RelayRequestKind::DeleteContent { content_id, @@ -381,7 +383,7 @@ impl StateNode { timestamp, ) .await - .map(|_| ()) + .map(|_| RelayOutcome::Done) } RelayRequestKind::InvalidateTokens { content_id, @@ -398,12 +400,49 @@ impl StateNode { timestamp, ) .await - .map(|_| ()) + .map(|_| RelayOutcome::Done) + } + RelayRequestKind::ReadContent { + content_id, + version, + auth_token, + request_signature, + timestamp, + } => { + let token = AuthToken::new(auth_token); + service_for_relay + .read_content_via_relay( + &content_id, + version.as_deref(), + &token, + Some(&request_signature), + timestamp, + ) + .await + .map(|(data, version)| RelayOutcome::Data { + data, + version, + }) + } + RelayRequestKind::ReadHistory { + content_id, + auth_token, + request_signature, + timestamp, + } => { + let token = AuthToken::new(auth_token); + service_for_relay + .read_history_via_relay( + &content_id, + &token, + Some(&request_signature), + timestamp, + ) + .await + .map(|versions| RelayOutcome::History { versions }) } }; - let _ = req - .reply - .send(result.map_err(|e| anyhow::anyhow!(e.to_string()))); + let _ = req.reply.send(result); } } } @@ -431,14 +470,15 @@ impl StateNode { match result { Ok(received) => { tracing::debug!( - "Received event from {}: {:?}", + "Received event from {:?}: {:?}", received.source, received.event.event_type() ); - // Forward to service for processing (with source PeerID for verification) + // Forward to service for processing, with the + // authenticated publisher PeerID for origin checks. match service - .handle_sync_event(&received.event, Some(&received.source)) + .handle_sync_event(&received.event, received.source.as_deref()) .await { Ok(outcome) => { diff --git a/monas-state-node/src/application_service/state_node_service.rs b/monas-state-node/src/application_service/state_node_service.rs index 882ca4e..9ad201b 100644 --- a/monas-state-node/src/application_service/state_node_service.rs +++ b/monas-state-node/src/application_service/state_node_service.rs @@ -18,11 +18,12 @@ use crate::port::authentication_service::AuthenticationService; use crate::port::authorization_service::{AuthorizationRequest, AuthorizationService}; use crate::port::content_repository::ContentRepository; use crate::port::event_publisher::EventPublisher; -use crate::port::peer_network::PeerNetwork; +use crate::port::peer_network::{PeerNetwork, RelayReadError, RelayReadErrorKind}; use crate::port::persistence::{ PersistentAccessControlRepository, PersistentContentRepository, PersistentNodeRegistry, }; use anyhow::Result; +use std::ops::ControlFlow; use std::sync::Arc; /// Result of applying an event. @@ -94,6 +95,105 @@ where max_add_member_count: usize, } +/// Where a relay candidate list came from, and therefore how much it can be +/// trusted. +/// +/// This distinction matters because a relay interprets what comes back from +/// whoever is in the list. Trusting a *negative* answer is only safe from a +/// peer we have some reason to believe actually holds the content. +/// +/// It does not affect whether credentials may be forwarded — see +/// [`ResolvedMembers::as_slice`] for why that is safe regardless of +/// provenance. +#[derive(Debug, Clone, Copy, PartialEq, Eq)] +enum MemberProvenance { + /// The peers come from a local `ContentNetwork` record, built from a + /// `ContentCreated` / `ContentNetworkManagerAdded` event that named this + /// node as a member. + /// + /// This is **not** a cryptographic attestation by the content owner. + /// Events carry no owner signature, so the membership set inside the event + /// is the publisher's own claim. What is verified is the *publisher*: + /// gossipsub runs in `Signed` + `Strict` mode, so the author field is + /// authenticated, and `handle_sync_event` binds it — `ContentCreated` must + /// come from the creator it names, and a membership change must come from + /// a node already in the network it changes. + /// + /// Two gaps remain, both requiring a protocol change to close (tracked + /// separately): the very first record for a content is accepted without a + /// prior membership to check against, and an authenticated member can + /// still claim a member set of its choosing. + /// + /// It is nevertheless a far better basis than [`Self::DhtGuess`], where + /// nothing at all ties a peer to the content — so a 401/403 from a listed + /// member is treated as a real verdict. + LocalRecord, + /// The peers are just DHT neighbours of `sha256(content_id)`. Nothing ties + /// them to this content: anyone able to place a Peer ID near the key lands + /// in this list, so a denial from one of them carries no authority. + DhtGuess, +} + +/// Relay candidates plus the provenance of the list. +struct ResolvedMembers { + members: Vec, + /// Kept for diagnostics and for #63: once membership is owner-signed, + /// `auth_verdict_is_authoritative` reads this again to restore the early + /// exit on a denial. Nothing branches on it today — see that method. + #[allow(dead_code)] + provenance: MemberProvenance, +} + +impl ResolvedMembers { + /// The relay candidates. + /// + /// Every candidate is handed the caller's token and request signature, and + /// that is fine even for an unproven peer: the token is either a + /// self-contained key id (a public key) or a delegated JWT whose audience + /// is likewise a public key, and authorization is proof-of-possession — the + /// request signature is verified against the JWT's `aud` key. A peer that + /// captures both cannot mint a new request, because it does not hold that + /// private key, and the signature it did capture is bound to one operation, + /// resource, body and timestamp. Mutations are single-use on top of that. + /// + /// So the candidate list is not truncated for unproven peers. Doing so + /// would cut failover to legitimate members — a real availability cost — + /// in exchange for no confidentiality gain, since a peer that sits first in + /// DHT distance order receives the credentials regardless of any cap. + fn as_slice(&self) -> &[String] { + &self.members + } + + /// Whether a negative authorization verdict from these peers may be treated + /// as final — i.e. may end the failover loop early. + /// + /// **Currently always false.** A denial is still kept and returned if no + /// candidate produces anything better; what this disables is *stopping* at + /// the first one. + /// + /// For [`MemberProvenance::DhtGuess`] the reason is direct: a single + /// hostile node squatting near the DHT key could otherwise deny every read + /// by answering 403 first, and even an honest but partially-synced replica + /// can answer 403 from a policy it has not finished replicating. + /// + /// [`MemberProvenance::LocalRecord`] used to return true here, on the + /// grounds that a listed member evaluated the caller against the real + /// policy. That does not hold while the record itself can be planted: the + /// first record for a content is accepted with no prior membership to check + /// against, so an attacker who wins that race lands in the list and its 403 + /// would end the loop — a permanent denial of service against a legitimate + /// caller, which is exactly what the DhtGuess case is guarding against. + /// + /// Restoring the early exit needs owner-signed membership (#63). Until + /// then, the cost of always continuing is bounded: one extra round trip per + /// remaining candidate on a genuine denial, against an availability attack + /// that is otherwise unbounded. Erring toward availability is the right + /// direction — the caller is refused either way, just later. + fn auth_verdict_is_authoritative(&self) -> bool { + false + } +} + /// No-op access control repository for backward compatibility. pub struct NoOpAccessControlRepository; @@ -222,12 +322,18 @@ where /// Authenticate a caller for read operations. /// + /// The request signature is bound to the specific `content_id` (message + /// `read:{content_id}:{timestamp}`), so a signature captured by one node + /// cannot be replayed to read other content. Mirrors the delete path, + /// which already signs over the content id. + /// /// Returns the authenticated identity on success. pub async fn authenticate_for_read( &self, token: &AuthToken, request_signature: Option<&[u8]>, timestamp: Option, + content_id: &str, ) -> Result { let auth_service = self.auth_service.as_ref().ok_or_else(|| { StateNodeError::InvalidConfiguration("Authentication not configured".to_string()) @@ -249,7 +355,7 @@ where token, sig, "read", - "content", + content_id, timestamp, None, ) @@ -411,7 +517,7 @@ where /// since they indicate a real permission problem rather than a connectivity issue. async fn relay_with_failover( &self, - members: &[String], + members: &ResolvedMembers, operation_name: &str, relay_fn: F, ) -> Result<(), StateNodeError> @@ -419,7 +525,11 @@ where F: Fn(String) -> Fut, Fut: std::future::Future>, { - for (i, member) in members.iter().enumerate() { + let authoritative = members.auth_verdict_is_authoritative(); + let candidates = members.as_slice(); + let mut auth_error: Option = None; + + for (i, member) in candidates.iter().enumerate() { match relay_fn(member.clone()).await { Ok(true) => return Ok(()), Ok(false) => { @@ -428,29 +538,51 @@ where operation_name, member, i + 1, - members.len() + candidates.len() ); } Err(e) => { let err_msg = e.to_string(); - // Don't failover on auth errors — they'll fail on every member if err_msg.contains("Authorization failed") || err_msg.contains("Authentication failed") { - return Err(Self::classify_relay_error(err_msg)); + // An attested member evaluated the request against the + // real policy, so every other member would answer the + // same — stop here. + if authoritative { + return Err(Self::classify_relay_error(err_msg)); + } + // Unproven candidate: its verdict proves nothing (it may + // not hold the content at all), so it must not end the + // loop — otherwise one hostile peer near the DHT key + // could block every write with a single fabricated 403. + // Remember it as the answer of last resort and continue. + tracing::warn!( + "Relay {} got an auth verdict from unproven candidate {} ({}/{}): {} \ + — continuing failover", + operation_name, + member, + i + 1, + candidates.len(), + err_msg + ); + auth_error.get_or_insert_with(|| Self::classify_relay_error(err_msg)); + continue; } tracing::warn!( "Relay {} failed on member {} ({}/{}): {}", operation_name, member, i + 1, - members.len(), + candidates.len(), err_msg ); } } } - Err(StateNodeError::NoAvailableMembers) + // Nothing succeeded. An auth verdict, even from an unproven candidate, + // is a more useful answer than "no members available". + Err(auth_error.unwrap_or(StateNodeError::NoAvailableMembers)) } /// Resolve the member nodes that own a content, for relay/sync purposes. @@ -465,7 +597,7 @@ where /// /// The local node is always excluded from the result. Returns /// `NoAvailableMembers` if no candidate remains. - async fn resolve_members(&self, content_id: &str) -> Result, StateNodeError> { + async fn resolve_members(&self, content_id: &str) -> Result { let local_record = self .content_repo .read() @@ -474,19 +606,24 @@ where .await .map_err(|e| StateNodeError::StorageError(e.to_string()))?; - let mut members: Vec = match local_record { - Some(network) => network.member_nodes_as_strings(), + let (mut members, provenance) = match local_record { + Some(network) => ( + network.member_nodes_as_strings(), + MemberProvenance::LocalRecord, + ), None => { // No local record: discover the members via the DHT. Request // `k + 1` candidates so that excluding ourselves still leaves // up to `k`, mirroring create_content's placement. let key = compute_dht_key(content_id); - self.peer_network + let peers = self + .peer_network .find_closest_peers(key, self.min_replication_factor + 1) .await .map_err(|e| { StateNodeError::NetworkError(NetworkError::ConnectionFailed(e.to_string())) - })? + })?; + (peers, MemberProvenance::DhtGuess) } }; @@ -494,76 +631,336 @@ where if members.is_empty() { return Err(StateNodeError::NoAvailableMembers); } - Ok(members) - } - - /// Ensure the content's CRDT state is available locally, fetching it from - /// member nodes if necessary (bug #93, read side). - /// - /// Read endpoints (data / history / version) read straight from the local - /// `crdt_repo`. A node that holds no local state for a content — e.g. a - /// client pointed its gateway at a node that is neither the creator nor a - /// member — would otherwise 404. This pulls the operations from a member - /// (discovered via `resolve_members`) and applies them locally so the - /// existing read path works unchanged. Applying the operations also brings - /// in the access policy, so the subsequent `verify_read_access` check is - /// evaluated against the real policy rather than an empty one. - /// - /// No-op when we already have local history. Returns `ContentNotFound` if - /// no member could supply the operations. - /// - /// SECURITY NOTE: this reuses the unauthenticated `fetch_operations` RPC, - /// so a non-member can pull any content's operations. Acceptable for the - /// single-user demo; a hardened version should add a read-relay RPC with - /// member-side authorization. Tracked in bug #93 follow-up. - pub async fn ensure_content_local(&self, content_id: &str) -> Result<(), StateNodeError> { - // Fast path: we already hold local history for this content. - let has_local = self + Ok(ResolvedMembers { + members, + provenance, + }) + } + + /// Whether this node actually replicates the content (holds its genesis + /// node in the local DAG). + /// + /// NOTE: `get_history` cannot be used for this — crsl-lib's + /// `linear_history` returns `[genesis]` even when no node exists (phantom + /// history), which would make such a check always pass. + pub async fn has_local_content(&self, content_id: &str) -> bool { + self.crdt_repo + .has_genesis(content_id) + .await + .unwrap_or(false) + } + + /// Authorize a read against the content's access policy (bug #93 hardened + /// read path). + /// + /// Authenticates the caller (token + request signature) and grants access + /// only when the caller is the content owner or the authorization service + /// grants `ReadContent`. + /// + /// Fail-closed in both directions: an error loading the policy denies, and + /// **a missing policy also denies**. A replica can legitimately hold a + /// genesis without its owner policy — the create operation carries + /// `access_policy: None` and the owner policy arrives as a separate + /// operation, and `apply_operations` tolerates partial application — so + /// treating "no policy" as public would expose ciphertext and history of + /// such content to any authenticated caller. + pub async fn authorize_read( + &self, + token: &AuthToken, + request_signature: Option<&[u8]>, + timestamp: Option, + content_id: &str, + ) -> Result<(), StateNodeError> { + let identity = self + .authenticate_for_read(token, request_signature, timestamp, content_id) + .await?; + + // Fail closed: an error loading the policy must deny, not fall through + // to the "no policy" allow below. + let policy = self .crdt_repo - .get_history(content_id) + .get_access_policy(content_id) .await - .map(|h| !h.is_empty()) - .unwrap_or(false); - if has_local { - return Ok(()); + .map_err(|e| { + StateNodeError::StorageError(format!( + "failed to load access policy for {}: {}", + content_id, e + )) + })?; + + if let Some(policy) = policy { + if policy.is_owner(&identity) { + return Ok(()); + } + + let authz_request = crate::port::authorization_service::AuthorizationRequest { + identity, + resource: ContentId::new(content_id.to_string())?, + capability: crate::domain::auth_capability::AuthCapability::ReadContent, + token: Some(token.clone()), + request_signature: request_signature.map(|s| s.to_vec()), + }; + if let Some(authz_service) = self.authz_service.as_ref() { + if let Ok(result) = authz_service.authorize(&authz_request).await { + if result.is_granted() { + return Ok(()); + } + } + } + + return Err(StateNodeError::AuthorizationFailed( + "Insufficient permissions: read access required".to_string(), + )); } - // Discover the members (local record, or DHT fallback) and pull ops. - let members = self.resolve_members(content_id).await?; + // No policy on this replica. This is not proof that the content is + // public — it is indistinguishable from "the owner policy operation has + // not been applied here (yet)". Deny rather than serve ciphertext and + // history without an authorization contract. + Err(StateNodeError::AuthorizationFailed(format!( + "no access policy is available for {content_id} on this node: refusing the read \ + (the policy may not have replicated here yet — retry, or read from a node that \ + has it)" + ))) + } + + /// Serve a relayed data read on a member node (bug #93 hardened read path). + /// + /// Called by the relay handler when a non-member node forwards a read. + /// Re-authenticates the original caller before touching the repository. + /// `version: None` serves the latest version. Returns `(data, version)`. + pub async fn read_content_via_relay( + &self, + content_id: &str, + version: Option<&str>, + token: &AuthToken, + request_signature: Option<&[u8]>, + timestamp: Option, + ) -> Result<(Vec, String), StateNodeError> { + self.authorize_read(token, request_signature, timestamp, content_id) + .await?; + let content_id_vo = ContentId::new(content_id.to_string())?; + match version { + Some(v) => { + let data = self + .crdt_repo + .get_version(content_id, v) + .await + .map_err(|e| StateNodeError::StorageError(e.to_string()))? + .ok_or(StateNodeError::ContentNotFound(content_id_vo))?; + Ok((data, v.to_string())) + } + None => self + .crdt_repo + .get_latest_with_version(content_id) + .await + .map_err(|e| StateNodeError::StorageError(e.to_string()))? + .ok_or(StateNodeError::ContentNotFound(content_id_vo)), + } + } + + /// Serve a relayed history read on a member node (bug #93 hardened read + /// path). Same authentication contract as `read_content_via_relay`. + pub async fn read_history_via_relay( + &self, + content_id: &str, + token: &AuthToken, + request_signature: Option<&[u8]>, + timestamp: Option, + ) -> Result, StateNodeError> { + self.authorize_read(token, request_signature, timestamp, content_id) + .await?; + + // Guard against crsl-lib's phantom `[genesis]` history for content we + // don't actually hold. + if !self.has_local_content(content_id).await { + return Err(StateNodeError::ContentNotFound(ContentId::new( + content_id.to_string(), + )?)); + } + + self.crdt_repo + .get_history(content_id) + .await + .map_err(|e| StateNodeError::StorageError(e.to_string())) + } + + /// Relay a data read to the content's members (bug #93 hardened read path, + /// caller side). + /// + /// Used by read endpoints on a node that does not replicate the content. + /// The caller's auth material is forwarded verbatim so the member can + /// re-authenticate; operations are never pulled to this node. + pub async fn relay_read_data( + &self, + content_id: &str, + version: Option<&str>, + token: &AuthToken, + request_signature: Option<&[u8]>, + timestamp: Option, + ) -> Result<(Vec, String), StateNodeError> { + let members = self.resolve_members(content_id).await?; + let authoritative = members.auth_verdict_is_authoritative(); + let sig: &[u8] = request_signature.unwrap_or(&[]); - for member in &members { + let mut best: Option = None; + for member in members.as_slice() { match self .peer_network - .fetch_operations(member, content_id, None) + .relay_read_content(member, content_id, version, token.as_str(), sig, timestamp) .await { - Ok(ops) if !ops.is_empty() => match self.crdt_repo.apply_operations(&ops).await { - Ok(_) => return Ok(()), - Err(e) => { - tracing::warn!( - "ensure_content_local: failed to apply ops from {} for {}: {}", - member, - content_id, - e - ); + Ok(result) => return Ok(result), + Err(e) => { + tracing::warn!( + "relay_read_data: member {} failed for {}: {}", + member, + content_id, + e + ); + match Self::record_relay_read_error(&mut best, e, authoritative) { + // An auth verdict from an attested member is + // authoritative — it DID evaluate the request against + // the real policy. Do not let a later member's + // transport failure overwrite it. + ControlFlow::Break(final_err) => { + return Err(Self::relay_read_error_to_state_error( + content_id, final_err, + )) + } + ControlFlow::Continue(()) => {} } - }, - Ok(_) => { - // Member returned no operations; try the next one. } + } + } + + Err(Self::relay_read_error_to_state_error( + content_id, + best.unwrap_or_else(|| RelayReadError::other("no members responded")), + )) + } + + /// Relay a history read to the content's members (bug #93 hardened read + /// path, caller side). + pub async fn relay_read_history( + &self, + content_id: &str, + token: &AuthToken, + request_signature: Option<&[u8]>, + timestamp: Option, + ) -> Result, StateNodeError> { + let members = self.resolve_members(content_id).await?; + let authoritative = members.auth_verdict_is_authoritative(); + let sig: &[u8] = request_signature.unwrap_or(&[]); + + let mut best: Option = None; + for member in members.as_slice() { + match self + .peer_network + .relay_read_history(member, content_id, token.as_str(), sig, timestamp) + .await + { + Ok(versions) => return Ok(versions), Err(e) => { tracing::warn!( - "ensure_content_local: failed to fetch ops from {} for {}: {}", + "relay_read_history: member {} failed for {}: {}", member, content_id, e ); + match Self::record_relay_read_error(&mut best, e, authoritative) { + ControlFlow::Break(final_err) => { + return Err(Self::relay_read_error_to_state_error( + content_id, final_err, + )) + } + ControlFlow::Continue(()) => {} + } + } + } + } + + Err(Self::relay_read_error_to_state_error( + content_id, + best.unwrap_or_else(|| RelayReadError::other("no members responded")), + )) + } + + /// Fold one member's relayed-read failure into the running best error. + /// + /// Auth verdicts (401/403) short-circuit the member loop **only when the + /// candidate list is attested** (`authoritative`): such a member evaluated + /// the caller against the real policy, so asking further members cannot + /// change the answer, and continuing would just leak that the content + /// exists to a caller who was already refused. + /// + /// When the list is a DHT guess, a 401/403 proves nothing — the responder + /// may not hold the content at all. Treating it as final would let one + /// hostile peer near the DHT key deny every read by answering 403 first, + /// and would also let an honest replica that has not finished replicating + /// the policy stop failover to a healthy one. So the verdict is remembered + /// (it is a better answer than a transport error, and is what the caller + /// sees if nothing better turns up) but the loop keeps going. + /// + /// NotFound is kept over transport errors but never short-circuits — a + /// lagging member may miss a version another member can still serve. + fn record_relay_read_error( + best: &mut Option, + err: RelayReadError, + authoritative: bool, + ) -> ControlFlow { + match err.kind { + RelayReadErrorKind::AuthenticationFailed | RelayReadErrorKind::AuthorizationFailed => { + if authoritative { + return ControlFlow::Break(err); + } + // Unproven responder: keep the verdict as the best answer so + // far (it beats NotFound and transport errors), but let the + // remaining candidates have their say. + *best = Some(err); + ControlFlow::Continue(()) + } + RelayReadErrorKind::NotFound => { + // Do not let a NotFound overwrite an auth verdict we are + // holding on to from an unproven peer: the verdict is the more + // specific answer. + if !matches!( + best.as_ref().map(|b| b.kind), + Some(RelayReadErrorKind::AuthenticationFailed) + | Some(RelayReadErrorKind::AuthorizationFailed) + ) { + *best = Some(err); } + ControlFlow::Continue(()) + } + RelayReadErrorKind::Other => { + if best.is_none() { + *best = Some(err); + } + ControlFlow::Continue(()) } } + } - Err(StateNodeError::ContentNotFound(content_id_vo)) + /// Convert a member's typed relayed-read verdict into the service error + /// the HTTP layer maps to a status code. + fn relay_read_error_to_state_error(content_id: &str, err: RelayReadError) -> StateNodeError { + match err.kind { + RelayReadErrorKind::NotFound => match ContentId::new(content_id.to_string()) { + Ok(cid) => StateNodeError::ContentNotFound(cid), + Err(_) => StateNodeError::StorageError(err.message), + }, + RelayReadErrorKind::AuthenticationFailed => { + StateNodeError::AuthenticationFailed(err.message) + } + RelayReadErrorKind::AuthorizationFailed => { + StateNodeError::AuthorizationFailed(err.message) + } + RelayReadErrorKind::Other => { + StateNodeError::NetworkError(NetworkError::ConnectionFailed(err.message)) + } + } } /// Register a new node. @@ -1751,10 +2148,56 @@ where Ok(()) } + /// Verify that the publisher of a membership-changing event is itself a + /// member of the network it is changing. + /// + /// Returns `Ok(())` when we hold no local record for the content: there is + /// no membership to check against, and no existing record to overwrite. + /// Also returns `Ok(())` when `source_peer_id` is `None` (the event did not + /// arrive over an authenticated channel — see `handle_sync_event`). + async fn verify_source_is_existing_member( + &self, + source_peer_id: Option<&str>, + content_id: &str, + ) -> Result<(), StateNodeError> { + let Some(source) = source_peer_id else { + return Ok(()); + }; + + let existing = self + .content_repo + .read() + .await + .get_content_network(content_id) + .await + .map_err(|e| StateNodeError::StorageError(e.to_string()))?; + + match existing { + Some(network) if !network.has_member_str(source) => { + tracing::warn!( + "Rejecting membership change for {} from non-member {}", + content_id, + source + ); + Err(StateNodeError::Internal(format!( + "Publisher {} is not a member of content network {}", + source, content_id + ))) + } + _ => Ok(()), + } + } + /// Handle a sync event from another node. /// - /// The `source_peer_id` parameter is used to verify that events claiming - /// to be from a particular node actually came from that peer's PeerID. + /// The `source_peer_id` parameter is the **authenticated publisher** of the + /// event (gossipsub `Message::source`, not the forwarding peer), used to + /// verify that events claiming to be from a particular node actually came + /// from that peer. + /// + /// `None` means the origin could not be established; origin-bound checks + /// are skipped in that case, so callers must pass the authenticated value + /// whenever one exists. /// /// Returns `ApplyOutcome::NeedsSync` when the caller should perform content /// synchronization (e.g., call `ContentSyncService::sync_from_peers`). @@ -1818,6 +2261,18 @@ where return Ok(ApplyOutcome::Ignored); } + // Unlike `ContentCreated`, this event names no publisher — + // `added_node_id` is the node being added, which is usually us. + // The publisher is whichever node ran the redundancy check, so + // the check that fits is membership: only a node already in the + // network we remember may change that network's member set. + // + // If we hold no record, there is nothing to check against and + // nothing being overwritten, so the event is accepted as + // bootstrap — the same position `ContentCreated` is in. + self.verify_source_is_existing_member(source_peer_id, content_id) + .await?; + // When handling sync events, we create network with NodeIds directly let content_id_vo = ContentId::new(content_id.clone())?; @@ -1848,6 +2303,14 @@ where removed_node_id, .. } => { + // Removal deletes or rewrites the record that decides whether a + // relay treats a peer's 403 as final, so the publisher must be + // a member of the network it is changing — otherwise any peer + // could evict us from our own record just by naming us as + // `removed_node_id`. Same rule as the Added arm. + self.verify_source_is_existing_member(source_peer_id, content_id) + .await?; + // If we were removed, delete the local network metadata if removed_node_id == &self.local_node_id { tracing::info!( @@ -1892,9 +2355,17 @@ where Event::ContentCreated { content_id, + creator_node_id, member_nodes, .. } => { + // The event asserts a membership set that this node will store + // and later act on, so the publisher must at least be the + // creator it claims to be. `create_content` always publishes + // with `creator_node_id == self.local_node_id`, so this holds + // for every legitimate event. + Self::verify_source_peer_id(source_peer_id, creator_node_id)?; + // Only store network metadata if we're a member if !member_nodes.contains(&self.local_node_id) { return Ok(ApplyOutcome::Ignored); @@ -1952,6 +2423,13 @@ where // Verify source PeerID matches claimed node ID Self::verify_source_peer_id(source_peer_id, deleted_by_node_id)?; + // That alone only proves the publisher is who it says it is — + // it names *itself*, so any authenticated peer would satisfy + // it. Deleting our record requires being a member of the + // network being deleted. + self.verify_source_is_existing_member(source_peer_id, content_id) + .await?; + // Skip if we initiated the deletion if deleted_by_node_id == &self.local_node_id { return Ok(ApplyOutcome::Ignored); @@ -2756,28 +3234,17 @@ mod tests { ); } - fn sample_operation(genesis_cid: &str) -> crate::port::content_repository::SerializedOperation { - crate::port::content_repository::SerializedOperation { - data: vec![0x01, 0x02], - genesis_cid: genesis_cid.to_string(), - author: "node-2".to_string(), - timestamp: 1, - node_timestamp: 1, - } - } - #[tokio::test] - async fn test_ensure_content_local_pulls_from_discovered_member() { - // Bug #93 (read side): a node with no local state for the content must - // pull the operations from a DHT-discovered member and apply them - // locally so the read endpoints work instead of 404ing. + async fn test_relay_read_data_maps_member_not_found() { + // Bug #93 (hardened read side): a node with no local state relays the + // read to a DHT-discovered member. When every member reports the + // content missing, the caller gets a typed ContentNotFound back. let node_registry = MockNodeRegistry::new(); let content_repo = Arc::new(RwLock::new(MockContentNetworkRepository::new())); let peer_network = Arc::new( MockPeerNetwork::new() .with_local_peer_id("node-1") - .with_closest_peers(vec!["node-2".to_string()]) - .with_fetched_operations(vec![sample_operation("content-1")]), + .with_closest_peers(vec!["node-2".to_string()]), ); let event_publisher = MockEventPublisher::new(); let crdt_repo = Arc::new(MockContentRepository::new()); @@ -2791,48 +3258,511 @@ mod tests { "node-1".to_string(), ); - let result = service.ensure_content_local("content-1").await; - assert!( - result.is_ok(), - "expected ops pull to succeed, got {result:?}" - ); - } - - #[tokio::test] - async fn test_ensure_content_local_errors_when_no_member_has_data() { - // No discoverable members → cannot pull → ContentNotFound. - let service = create_test_service("node-1"); - let result = service.ensure_content_local("content-1").await; - assert!(result.is_err()); - assert!(result - .unwrap_err() - .to_string() - .contains("No available member nodes")); + let result = service + .relay_read_data( + "content-1", + None, + &test_token(), + Some(&test_request_signature()), + None, + ) + .await; + assert!(matches!(result, Err(StateNodeError::ContentNotFound(_)))); } #[tokio::test] - async fn test_handle_sync_event_node_created() { - let service = create_test_service("node-1"); - - let event = Event::NodeCreated { - node_id: "node-2".to_string(), - total_capacity: 2000, - available_capacity: 1500, - timestamp: 12345, - }; + async fn test_relay_read_data_returns_member_payload() { + // Happy path: the member serves the read and the payload comes back. + let node_registry = MockNodeRegistry::new(); + let content_repo = Arc::new(RwLock::new(MockContentNetworkRepository::new())); + let peer_network = Arc::new( + MockPeerNetwork::new() + .with_local_peer_id("node-1") + .with_closest_peers(vec!["node-2".to_string()]) + .with_relay_read_data(b"cipher".to_vec(), "v1"), + ); + let event_publisher = MockEventPublisher::new(); + let crdt_repo = Arc::new(MockContentRepository::new()); - let outcome = service.handle_sync_event(&event, None).await.unwrap(); - assert_eq!(outcome, ApplyOutcome::Applied); + let service: TestService = StateNodeService::new( + node_registry, + content_repo, + peer_network, + event_publisher, + crdt_repo, + "node-1".to_string(), + ); - // Verify node was stored - let stored = service.get_node("node-2").await.unwrap().unwrap(); - assert_eq!(stored.total_capacity, 2000); - assert_eq!(stored.available_capacity, 1500); + let result = service + .relay_read_data( + "content-1", + None, + &test_token(), + Some(&test_request_signature()), + None, + ) + .await + .expect("relayed read should succeed"); + assert_eq!(result, (b"cipher".to_vec(), "v1".to_string())); } #[tokio::test] - async fn test_handle_sync_event_content_created_as_member() { - let service = create_test_service("node-1"); + async fn test_relay_read_data_returns_member_auth_verdict() { + // A member's 403 verdict must come back typed (not as a generic + // network error), so the HTTP layer returns the member's decision. + use crate::port::peer_network::{RelayReadError, RelayReadErrorKind}; + let node_registry = MockNodeRegistry::new(); + let content_repo = Arc::new(RwLock::new(MockContentNetworkRepository::new())); + let peer_network = Arc::new( + MockPeerNetwork::new() + .with_local_peer_id("node-1") + .with_closest_peers(vec!["node-2".to_string(), "node-3".to_string()]) + .with_relay_read_error(RelayReadError { + kind: RelayReadErrorKind::AuthorizationFailed, + message: "Insufficient permissions: read access required".to_string(), + }), + ); + let event_publisher = MockEventPublisher::new(); + let crdt_repo = Arc::new(MockContentRepository::new()); + + let service: TestService = StateNodeService::new( + node_registry, + content_repo, + peer_network, + event_publisher, + crdt_repo, + "node-1".to_string(), + ); + + let result = service + .relay_read_data( + "content-1", + None, + &test_token(), + Some(&test_request_signature()), + None, + ) + .await; + assert!(matches!( + result, + Err(StateNodeError::AuthorizationFailed(_)) + )); + } + + /// 未証明の DHT 候補が返した 401/403 で member loop を打ち切ってはいけない。 + /// 打ち切ると、DHT キーの近くに Peer ID を置いた 1 台が 403 を返すだけで + /// あらゆる read を止められる(可用性への攻撃)。正規 member でも policy の + /// 複製が終わっていなければ 403 を返し得るので、健全なレプリカへの failover + /// を潰さないためにも継続する必要がある。 + #[test] + fn unproven_peer_auth_verdict_does_not_stop_failover() { + use crate::port::peer_network::{RelayReadError, RelayReadErrorKind}; + + let mut best = None; + let flow = StateNodeService::< + MockNodeRegistry, + MockContentNetworkRepository, + MockPeerNetwork, + MockEventPublisher, + MockContentRepository, + >::record_relay_read_error( + &mut best, + RelayReadError { + kind: RelayReadErrorKind::AuthorizationFailed, + message: "fabricated 403".to_string(), + }, + false, // 未証明の候補 + ); + assert!( + matches!(flow, ControlFlow::Continue(())), + "an unproven peer's verdict must not end the loop" + ); + // ただし答えとしては保持する(次の候補が何も返さなければこれを返す) + assert!(matches!( + best.as_ref().map(|b| b.kind), + Some(RelayReadErrorKind::AuthorizationFailed) + )); + + // 後続候補の NotFound で auth verdict を上書きしない。 + // 上書きすると「拒否された」が「存在しない」に化ける。 + let flow = StateNodeService::< + MockNodeRegistry, + MockContentNetworkRepository, + MockPeerNetwork, + MockEventPublisher, + MockContentRepository, + >::record_relay_read_error( + &mut best, + RelayReadError { + kind: RelayReadErrorKind::NotFound, + message: "not here".to_string(), + }, + false, + ); + assert!(matches!(flow, ControlFlow::Continue(()))); + assert!(matches!( + best.as_ref().map(|b| b.kind), + Some(RelayReadErrorKind::AuthorizationFailed) + )); + } + + /// `record_relay_read_error` は「権威あり」と言われれば打ち切る。 + /// + /// ただし現在この `true` を渡す呼び出し側は無い + /// ([`ResolvedMembers::auth_verdict_is_authoritative`] は常に false)。 + /// レコード自体が最初の 1 通で植え付けられる間は、そこに載った peer の + /// 403 も最終判断にはできないためである。owner 署名付き membership + /// (#63)が入れば早期打ち切りを戻せるので、その配線だけは残してある。 + #[test] + fn record_relay_read_error_breaks_when_told_the_verdict_is_authoritative() { + use crate::port::peer_network::{RelayReadError, RelayReadErrorKind}; + + let mut best = None; + let flow = StateNodeService::< + MockNodeRegistry, + MockContentNetworkRepository, + MockPeerNetwork, + MockEventPublisher, + MockContentRepository, + >::record_relay_read_error( + &mut best, + RelayReadError { + kind: RelayReadErrorKind::AuthorizationFailed, + message: "real verdict".to_string(), + }, + true, // attested member + ); + assert!( + matches!(flow, ControlFlow::Break(_)), + "an attested member's verdict is authoritative" + ); + } + + /// 出自は verdict の扱いだけを変え、候補リストそのものは削らない。 + /// + /// credential は未証明の相手へ渡っても構わない。token は自己完結型 key id + /// (公開鍵)か委譲 JWT で、その JWT の `aud` もまた公開鍵であり、認可は + /// proof-of-possession だからである — リクエスト署名は `aud` の鍵に対して + /// 検証されるので、両方を傍受した相手も秘密鍵を持たない以上、新しい + /// リクエストを作れない。傍受した署名自体も操作・リソース・body・timestamp + /// に束縛され、mutation はさらに使い切りである。 + /// + /// 逆に候補を削ると、正当な member への failover が減って可用性だけが + /// 落ちる。DHT 距離順で先頭に来る相手は上限があろうと credential を + /// 受け取るので、機密性は何も改善しない。 + #[test] + fn candidate_list_is_not_truncated_by_provenance() { + let many: Vec = (0..10).map(|i| format!("peer-{i}")).collect(); + + for provenance in [MemberProvenance::LocalRecord, MemberProvenance::DhtGuess] { + let resolved = ResolvedMembers { + members: many.clone(), + provenance, + }; + assert_eq!( + resolved.as_slice().len(), + many.len(), + "failover must reach every candidate ({provenance:?})" + ); + } + } + + #[tokio::test] + async fn test_authorize_read_allows_owner() { + let service = create_test_service("node-1"); + let owner = Identity::user("test-user".to_string()).unwrap(); + let policy = crate::domain::access_policy::AccessPolicy::new( + ContentId::new("content-1".to_string()).unwrap(), + owner, + ); + service + .crdt_repo + .access_policies + .lock() + .await + .insert("content-1".to_string(), policy); + + let result = service + .authorize_read( + &test_token(), + Some(&test_request_signature()), + None, + "content-1", + ) + .await; + assert!(result.is_ok(), "owner must be allowed: {result:?}"); + } + + /// policy が無いレプリカでの read は拒否する(fail-closed)。 + /// + /// create genesis は `access_policy: None` で作られ owner policy は別 + /// operation として届くため、「genesis はあるが policy が無い」状態は + /// 部分同期で実際に起こり得る。これを public 扱いにすると、認可契約の + /// 無いまま暗号文と履歴を任意の認証済み caller に渡してしまう。 + #[tokio::test] + async fn test_authorize_read_denies_when_policy_is_missing() { + let service = create_test_service("node-1"); + let result = service + .authorize_read( + &test_token(), + Some(&test_request_signature()), + None, + "content-1", + ) + .await; + match result { + Err(StateNodeError::AuthorizationFailed(msg)) => { + assert!(msg.contains("no access policy"), "msg={msg}"); + } + other => panic!("expected AuthorizationFailed, got: {other:?}"), + } + } + + /// 「genesis はレプリカにあるが owner policy がまだ届いていない」状態を + /// 直接再現し、非 owner の read が拒否されることを確認する。 + /// これが塞ぐ実シナリオ(create の Create payload は `access_policy: None` + /// で、owner policy は別 operation として届く)。 + #[tokio::test] + async fn test_authorize_read_denies_on_genesis_only_replica() { + let service = create_test_service("node-1"); + + // genesis(コンテンツ本体)だけが存在し、access_policies は空のまま + service + .crdt_repo + .contents + .lock() + .await + .insert("content-genesis-only".to_string(), b"ciphertext".to_vec()); + assert!( + service + .crdt_repo + .access_policies + .lock() + .await + .get("content-genesis-only") + .is_none(), + "precondition: replica must not have the owner policy yet" + ); + + let result = service + .authorize_read( + &test_token(), + Some(&test_request_signature()), + None, + "content-genesis-only", + ) + .await; + + assert!( + matches!(result, Err(StateNodeError::AuthorizationFailed(_))), + "genesis-only replica must not serve reads without a policy: {result:?}" + ); + } + + /// The read request signature must be verified against a message bound to + /// the specific content id (`read:{content_id}:{timestamp}`), not a + /// generic `read:content:{timestamp}`. This keeps a signature forwarded to + /// relay members (or leaked to a non-member node) from being replayed to + /// read other content (PR #54 review). + #[tokio::test] + async fn test_read_signature_message_is_bound_to_content_id() { + struct CapturingAuthService { + messages: Arc>>, + } + + #[async_trait::async_trait] + impl AuthenticationService for CapturingAuthService { + async fn authenticate( + &self, + token: &AuthToken, + _context: Option<&crate::port::auth_token::AuthContext>, + ) -> Result { + Identity::user(token.as_str().to_string()) + .map_err(|e| anyhow::anyhow!(e.to_string())) + } + + async fn is_valid(&self, token: &AuthToken) -> Result { + Ok(!token.is_empty()) + } + + async fn verify_request_signature( + &self, + _token: &AuthToken, + _signature: &[u8], + message: &str, + _timestamp: Option, + ) -> Result<()> { + self.messages.lock().unwrap().push(message.to_string()); + Ok(()) + } + + async fn verify_jwt_signature(&self, _token: &AuthToken) -> Result<()> { + Ok(()) + } + + async fn get_issuer(&self, token: &AuthToken) -> Result> { + Ok(Some( + Identity::user(token.as_str().to_string()) + .map_err(|e| anyhow::anyhow!(e.to_string()))?, + )) + } + } + + let messages = Arc::new(std::sync::Mutex::new(Vec::new())); + let node_registry = MockNodeRegistry::new(); + let content_repo = Arc::new(RwLock::new(MockContentNetworkRepository::new())); + let peer_network = Arc::new(MockPeerNetwork::new().with_local_peer_id("node-1")); + let event_publisher = MockEventPublisher::new(); + let crdt_repo = Arc::new(MockContentRepository::new()); + let service: TestService = StateNodeService::new( + node_registry, + content_repo, + peer_network, + event_publisher, + crdt_repo, + "node-1".to_string(), + ) + .with_authentication_service(CapturingAuthService { + messages: Arc::clone(&messages), + }); + + service + .authenticate_for_read( + &test_token(), + Some(&test_request_signature()), + Some(1234), + "content-abc", + ) + .await + .expect("authentication should succeed"); + + let captured = messages.lock().unwrap(); + assert_eq!( + captured.as_slice(), + ["read:content-abc:1234"], + "read signature message must include the content id" + ); + } + + #[tokio::test] + async fn test_authorize_read_denies_non_owner_when_authz_denies() { + struct DenyAllAuthorizationService; + #[async_trait::async_trait] + impl AuthorizationService for DenyAllAuthorizationService { + async fn authorize( + &self, + _request: &AuthorizationRequest, + ) -> Result { + Ok(AuthorizationResult::Denied { + reason: "no".to_string(), + }) + } + } + + let node_registry = MockNodeRegistry::new(); + let content_repo = Arc::new(RwLock::new(MockContentNetworkRepository::new())); + let peer_network = Arc::new(MockPeerNetwork::new().with_local_peer_id("node-1")); + let event_publisher = MockEventPublisher::new(); + let crdt_repo = Arc::new(MockContentRepository::new()); + + let service: TestService = StateNodeService::new( + node_registry, + content_repo, + peer_network, + event_publisher, + crdt_repo, + "node-1".to_string(), + ) + .with_authentication_service(TestAuthService) + .with_authorization_service(DenyAllAuthorizationService); + + let owner = Identity::user("someone-else".to_string()).unwrap(); + let policy = crate::domain::access_policy::AccessPolicy::new( + ContentId::new("content-1".to_string()).unwrap(), + owner, + ); + service + .crdt_repo + .access_policies + .lock() + .await + .insert("content-1".to_string(), policy); + + let result = service + .authorize_read( + &test_token(), + Some(&test_request_signature()), + None, + "content-1", + ) + .await; + assert!(matches!( + result, + Err(StateNodeError::AuthorizationFailed(_)) + )); + } + + #[tokio::test] + async fn test_authorize_read_fails_closed_on_policy_error() { + // A policy-store failure must deny, not fall through to the + // "no policy -> allow" branch. + let service = create_test_service("node-1"); + *service.crdt_repo.access_policy_error.lock().await = true; + + let result = service + .authorize_read( + &test_token(), + Some(&test_request_signature()), + None, + "content-1", + ) + .await; + assert!(matches!(result, Err(StateNodeError::StorageError(_)))); + } + + #[tokio::test] + async fn test_relay_read_data_errors_when_no_members() { + // No discoverable members → nothing to relay to → NoAvailableMembers. + let service = create_test_service("node-1"); + let result = service + .relay_read_data( + "content-1", + None, + &test_token(), + Some(&test_request_signature()), + None, + ) + .await; + assert!(result.is_err()); + assert!(result + .unwrap_err() + .to_string() + .contains("No available member nodes")); + } + + #[tokio::test] + async fn test_handle_sync_event_node_created() { + let service = create_test_service("node-1"); + + let event = Event::NodeCreated { + node_id: "node-2".to_string(), + total_capacity: 2000, + available_capacity: 1500, + timestamp: 12345, + }; + + let outcome = service.handle_sync_event(&event, None).await.unwrap(); + assert_eq!(outcome, ApplyOutcome::Applied); + + // Verify node was stored + let stored = service.get_node("node-2").await.unwrap().unwrap(); + assert_eq!(stored.total_capacity, 2000); + assert_eq!(stored.available_capacity, 1500); + } + + #[tokio::test] + async fn test_handle_sync_event_content_created_as_member() { + let service = create_test_service("node-1"); let event = Event::ContentCreated { content_id: "content-1".to_string(), @@ -2861,6 +3791,203 @@ mod tests { assert!(network.has_member_str("node-2")); } + /// A peer cannot plant a `ContentNetwork` record by publishing a + /// `ContentCreated` that names someone else as the creator. + /// + /// Without this check any node could name us in `member_nodes` for a + /// content of its choosing, and the record would read back as + /// `MemberProvenance::LocalRecord` — which decides whether a 403 from a + /// listed peer is treated as final. + #[tokio::test] + async fn content_created_from_a_peer_other_than_the_creator_is_rejected() { + let service = create_test_service("node-1"); + + let event = Event::ContentCreated { + content_id: "content-1".to_string(), + creator_node_id: "node-2".to_string(), + content_size: 100, + member_nodes: vec!["node-1".to_string(), "node-2".to_string()], + timestamp: 12345, + }; + + // Published by node-9, which is neither the claimed creator nor in the + // member set. + let result = service.handle_sync_event(&event, Some("node-9")).await; + assert!(result.is_err(), "unrelated publisher must be rejected"); + + let network = service + .get_content_network_for_test("content-1") + .await + .unwrap(); + assert!(network.is_none(), "no record may be planted"); + } + + /// The legitimate path still works: `create_content` publishes with + /// `creator_node_id == local_node_id`, so the authenticated publisher + /// matches the claimed creator. + #[tokio::test] + async fn content_created_from_the_real_creator_is_accepted() { + let service = create_test_service("node-1"); + + let event = Event::ContentCreated { + content_id: "content-1".to_string(), + creator_node_id: "node-2".to_string(), + content_size: 100, + member_nodes: vec!["node-1".to_string(), "node-2".to_string()], + timestamp: 12345, + }; + + let outcome = service + .handle_sync_event(&event, Some("node-2")) + .await + .unwrap(); + assert_eq!( + outcome, + ApplyOutcome::NeedsSync { + content_id: "content-1".to_string() + } + ); + + let network = service + .get_content_network_for_test("content-1") + .await + .unwrap() + .expect("record stored"); + assert!(network.has_member_str("node-2")); + } + + /// A non-member cannot evict us from our own record by naming us as the + /// removed node. + /// + /// This arm had no publisher check at all, so a single event from any peer + /// deleted the `ContentNetwork` record — which is what decides whether a + /// relay treats a peer's 403 as final. + #[tokio::test] + async fn removal_from_a_non_member_cannot_delete_our_record() { + let content_repo = Arc::new(RwLock::new( + MockContentNetworkRepository::new() + .with_network(create_test_network("content-1", vec!["node-1", "node-2"])), + )); + let service: TestService = StateNodeService::new( + MockNodeRegistry::new(), + content_repo, + Arc::new(MockPeerNetwork::new().with_local_peer_id("node-1")), + MockEventPublisher::new(), + Arc::new(MockContentRepository::new()), + "node-1".to_string(), + ) + .with_authentication_service(TestAuthService) + .with_authorization_service(AllowAllAuthorizationService); + + let event = Event::ContentNetworkManagerRemoved { + content_id: "content-1".to_string(), + removed_node_id: "node-1".to_string(), + member_nodes: vec!["node-2".to_string()], + reason: "low_capacity".to_string(), + timestamp: 12345, + }; + + let result = service.handle_sync_event(&event, Some("node-9")).await; + assert!(result.is_err(), "non-member publisher must be rejected"); + + assert!( + service + .get_content_network_for_test("content-1") + .await + .unwrap() + .is_some(), + "our record must survive" + ); + } + + /// A non-member cannot delete our record with a `ContentDeleted` event. + /// + /// `verify_source_peer_id` alone does not help here: the event names its + /// own publisher, so any authenticated peer satisfies it. + #[tokio::test] + async fn content_deleted_from_a_non_member_cannot_delete_our_record() { + let content_repo = Arc::new(RwLock::new( + MockContentNetworkRepository::new() + .with_network(create_test_network("content-1", vec!["node-1", "node-2"])), + )); + let service: TestService = StateNodeService::new( + MockNodeRegistry::new(), + content_repo, + Arc::new(MockPeerNetwork::new().with_local_peer_id("node-1")), + MockEventPublisher::new(), + Arc::new(MockContentRepository::new()), + "node-1".to_string(), + ) + .with_authentication_service(TestAuthService) + .with_authorization_service(AllowAllAuthorizationService); + + // node-9 both publishes and names itself — the self-claim check passes. + let event = Event::ContentDeleted { + content_id: "content-1".to_string(), + deleted_by_node_id: "node-9".to_string(), + timestamp: 12345, + }; + + let result = service.handle_sync_event(&event, Some("node-9")).await; + assert!(result.is_err(), "non-member publisher must be rejected"); + + assert!( + service + .get_content_network_for_test("content-1") + .await + .unwrap() + .is_some(), + "our record must survive" + ); + } + + /// A non-member cannot rewrite the member set of a network we already hold. + #[tokio::test] + async fn membership_change_from_a_non_member_is_rejected() { + let node_registry = MockNodeRegistry::new(); + let content_repo = Arc::new(RwLock::new( + MockContentNetworkRepository::new() + .with_network(create_test_network("content-1", vec!["node-1", "node-2"])), + )); + let peer_network = Arc::new(MockPeerNetwork::new().with_local_peer_id("node-1")); + + let service: TestService = StateNodeService::new( + node_registry, + content_repo, + peer_network, + MockEventPublisher::new(), + Arc::new(MockContentRepository::new()), + "node-1".to_string(), + ) + .with_authentication_service(TestAuthService) + .with_authorization_service(AllowAllAuthorizationService); + + // node-9 is not in the network, but tries to add itself to it. + let event = Event::ContentNetworkManagerAdded { + content_id: "content-1".to_string(), + added_node_id: "node-9".to_string(), + member_nodes: vec![ + "node-1".to_string(), + "node-2".to_string(), + "node-9".to_string(), + ], + timestamp: 12345, + }; + + let result = service.handle_sync_event(&event, Some("node-9")).await; + assert!(result.is_err(), "non-member publisher must be rejected"); + + let network = service + .get_content_network_for_test("content-1") + .await + .unwrap() + .expect("record still present"); + assert!( + !network.has_member_str("node-9"), + "member set must be unchanged" + ); + } + #[tokio::test] async fn test_handle_sync_event_content_created_not_member() { let service = create_test_service("node-1"); diff --git a/monas-state-node/src/bin/test_auth_generator.rs b/monas-state-node/src/bin/test_auth_generator.rs index 3aaee5c..67fe9c4 100644 --- a/monas-state-node/src/bin/test_auth_generator.rs +++ b/monas-state-node/src/bin/test_auth_generator.rs @@ -90,8 +90,8 @@ fn generate_test_auth_data() { // Machine-readable output for script consumption println!("PRIVATE_KEY={}", hex::encode(private_key_bytes)); - println!("PUBLIC_KEY={}", public_key_hex); - println!("KEY_ID=user:{}", public_key_hex); + println!("PUBLIC_KEY={public_key_hex}"); + println!("KEY_ID=user:{public_key_hex}"); } /// Sign a request with the correct message format. diff --git a/monas-state-node/src/infrastructure/auth/ucan_adapter.rs b/monas-state-node/src/infrastructure/auth/ucan_adapter.rs index 2575dd8..42e84c0 100644 --- a/monas-state-node/src/infrastructure/auth/ucan_adapter.rs +++ b/monas-state-node/src/infrastructure/auth/ucan_adapter.rs @@ -485,7 +485,11 @@ mod tests { ) -> Result, String)>> { unimplemented!() } - async fn get_version(&self, _version_cid: &str) -> Result>> { + async fn get_version( + &self, + _genesis_cid: &str, + _version_cid: &str, + ) -> Result>> { unimplemented!() } async fn get_history(&self, _genesis_cid: &str) -> Result> { diff --git a/monas-state-node/src/infrastructure/crdt_repository.rs b/monas-state-node/src/infrastructure/crdt_repository.rs index ec5487f..cf04d83 100644 --- a/monas-state-node/src/infrastructure/crdt_repository.rs +++ b/monas-state-node/src/infrastructure/crdt_repository.rs @@ -214,11 +214,20 @@ impl ContentRepository for CrslCrdtRepository { } } - async fn get_version(&self, version_cid: &str) -> Result>> { + async fn get_version(&self, genesis_cid: &str, version_cid: &str) -> Result>> { + let genesis = Self::parse_cid(genesis_cid)?; let cid = Self::parse_cid(version_cid)?; let repo = self.repo.lock(); + // Scope the lookup to this content's DAG: read authorization is + // granted per content, so serving a version from a different series + // would be a cross-content read. Unknown nodes resolve to None. + match repo.get_genesis(&cid) { + Ok(g) if g == genesis => {} + _ => return Ok(None), + } + match repo.dag.get_node(&cid) { Ok(Some(node)) => Ok(Some(node.payload().data.clone())), Ok(None) => Ok(None), @@ -702,6 +711,37 @@ mod tests { .unwrap()); } + #[tokio::test] + async fn test_get_version_rejects_cross_content_lookup() { + // Read authorization is granted per content, so a version CID from a + // different series must not be readable through another genesis. + let tmp = tempdir().unwrap(); + let repo = CrslCrdtRepository::open(tmp.path()).unwrap(); + + let a = repo + .create_content(b"content A", "author", None) + .await + .unwrap(); + let b = repo + .create_content(b"content B", "author", None) + .await + .unwrap(); + + // Sanity: within the right series the version resolves. + assert!(repo + .get_version(&a.genesis_cid, &a.version_cid) + .await + .unwrap() + .is_some()); + + // Cross-content: B's version through A's genesis must be None. + assert!(repo + .get_version(&a.genesis_cid, &b.version_cid) + .await + .unwrap() + .is_none()); + } + #[tokio::test] async fn test_get_history() { let tmp = tempdir().unwrap(); @@ -730,7 +770,10 @@ mod tests { let data = b"Test content"; let result = repo.create_content(data, "author", None).await.unwrap(); - let retrieved = repo.get_version(&result.version_cid).await.unwrap(); + let retrieved = repo + .get_version(&result.genesis_cid, &result.version_cid) + .await + .unwrap(); assert_eq!(retrieved, Some(data.to_vec())); } diff --git a/monas-state-node/src/infrastructure/gossipsub_publisher.rs b/monas-state-node/src/infrastructure/gossipsub_publisher.rs index 4fadc9f..2d35ec7 100644 --- a/monas-state-node/src/infrastructure/gossipsub_publisher.rs +++ b/monas-state-node/src/infrastructure/gossipsub_publisher.rs @@ -136,6 +136,7 @@ impl EventPublisher for GossipsubEventPublisher

{ mod tests { use super::*; use crate::port::content_repository::SerializedOperation; + use crate::port::peer_network::{RelayReadError, RelayReadErrorKind}; use std::collections::HashMap; /// Mock PeerNetwork for testing. @@ -281,6 +282,35 @@ mod tests { Ok(true) } + async fn relay_read_content( + &self, + _peer_id: &str, + content_id: &str, + _version: Option<&str>, + _auth_token: &str, + _request_signature: &[u8], + _timestamp: Option, + ) -> std::result::Result<(Vec, String), RelayReadError> { + Err(RelayReadError { + kind: RelayReadErrorKind::NotFound, + message: format!("Content not found: {}", content_id), + }) + } + + async fn relay_read_history( + &self, + _peer_id: &str, + content_id: &str, + _auth_token: &str, + _request_signature: &[u8], + _timestamp: Option, + ) -> std::result::Result, RelayReadError> { + Err(RelayReadError { + kind: RelayReadErrorKind::NotFound, + message: format!("Content not found: {}", content_id), + }) + } + async fn connected_peer_count(&self) -> usize { 0 } diff --git a/monas-state-node/src/infrastructure/network/libp2p_network.rs b/monas-state-node/src/infrastructure/network/libp2p_network.rs index 4a91525..9e46859 100644 --- a/monas-state-node/src/infrastructure/network/libp2p_network.rs +++ b/monas-state-node/src/infrastructure/network/libp2p_network.rs @@ -15,6 +15,7 @@ use crate::domain::events::Event; use crate::infrastructure::disk_capacity; use crate::port::content_repository::{ContentRepository, SerializedOperation}; use crate::port::peer_network::PeerNetwork; +use crate::port::peer_network::{RelayReadError, RelayReadErrorKind}; use anyhow::{Context, Result}; use async_trait::async_trait; @@ -42,7 +43,18 @@ const PEER_NETWORK_TIMEOUT: Duration = Duration::from_secs(30); /// which processes them using StateNodeService. pub struct RelayRequest { pub kind: RelayRequestKind, - pub reply: oneshot::Sender>, + pub reply: + oneshot::Sender>, +} + +/// Result payload of a processed relay request. +pub enum RelayOutcome { + /// Write relays (update/delete/invalidate) complete without a payload. + Done, + /// Relayed data read: raw bytes plus the version that was served. + Data { data: Vec, version: String }, + /// Relayed history read. + History { versions: Vec }, } /// The kind of relay request. @@ -66,6 +78,20 @@ pub enum RelayRequestKind { request_signature: Vec, timestamp: Option, }, + ReadContent { + content_id: String, + /// `None` reads the latest version; `Some(v)` a specific version CID. + version: Option, + auth_token: String, + request_signature: Vec, + timestamp: Option, + }, + ReadHistory { + content_id: String, + auth_token: String, + request_signature: Vec, + timestamp: Option, + }, } /// Gossipsub message received from the network. @@ -82,8 +108,21 @@ pub struct GossipsubMessage { /// Parsed domain event received from Gossipsub. #[derive(Debug, Clone)] pub struct ReceivedEvent { - /// The source peer ID. - pub source: String, + /// The peer that **published** the message, as authenticated by gossipsub. + /// + /// This is `Message::source`, not `propagation_source`: the mesh forwards + /// messages, so the peer that handed us the bytes is generally not the one + /// that produced them. Under `MessageAuthenticity::Signed` + + /// `ValidationMode::Strict` the author field is required and the message + /// signature is verified against it before delivery, so a forwarder cannot + /// alter it. Authorization checks must use this field — using the + /// forwarder would both accept forged origins and reject honest multi-hop + /// delivery. + /// + /// `None` only if a message somehow arrives without an author, which + /// Strict mode rejects; callers treat it as "unverifiable" and skip + /// origin-bound checks rather than trusting it. + pub source: Option, /// The parsed domain event. pub event: Event, } @@ -210,6 +249,23 @@ enum SwarmCommand { timestamp: Option, reply: oneshot::Sender>, }, + RelayReadContent { + peer_id: PeerId, + content_id: String, + version: Option, + auth_token: String, + request_signature: Vec, + timestamp: Option, + reply: RelayReadReply, + }, + RelayReadHistory { + peer_id: PeerId, + content_id: String, + auth_token: String, + request_signature: Vec, + timestamp: Option, + reply: RelayHistoryReply, + }, /// Send a response back through a ResponseChannel. /// Used by spawned relay tasks to send responses without blocking the swarm loop. SendRelayResponse { @@ -221,6 +277,26 @@ enum SwarmCommand { /// TTL for pending requests. Entries older than this are cleaned up to prevent memory leaks. const PENDING_REQUEST_TTL: Duration = Duration::from_secs(120); +/// Reply payload of a relayed data read: `(data, served_version)`. +type RelayReadReply = oneshot::Sender, String), RelayReadError>>; +/// Reply payload of a relayed history read. +type RelayHistoryReply = oneshot::Sender, RelayReadError>>; + +/// Map a member-side service error to the wire verdict for relayed reads. +fn relay_read_error_kind(e: &crate::domain::errors::StateNodeError) -> RelayReadErrorKind { + use crate::domain::errors::StateNodeError as E; + match e { + E::ContentNotFound(_) => RelayReadErrorKind::NotFound, + E::AuthenticationFailed(_) | E::InvalidUcanToken(_) => { + RelayReadErrorKind::AuthenticationFailed + } + E::AuthorizationFailed(_) | E::PermissionDenied(_) => { + RelayReadErrorKind::AuthorizationFailed + } + _ => RelayReadErrorKind::Other, + } +} + /// Pending requests tracking with TTL support. /// /// Each request tracks its creation time. A periodic sweep removes entries @@ -238,6 +314,8 @@ struct PendingRequests { relay_update_queries: HashMap>>, relay_delete_queries: HashMap>>, relay_invalidate_tokens_queries: HashMap>>, + relay_read_queries: HashMap, + relay_history_queries: HashMap, /// Timestamps for all pending request IDs, used for TTL-based cleanup. timestamps: HashMap, } @@ -261,6 +339,8 @@ impl PendingRequests { self.relay_delete_queries.retain(|_, s| !s.is_closed()); self.relay_invalidate_tokens_queries .retain(|_, s| !s.is_closed()); + self.relay_read_queries.retain(|_, s| !s.is_closed()); + self.relay_history_queries.retain(|_, s| !s.is_closed()); // Clean up expired timestamps self.timestamps @@ -756,6 +836,46 @@ impl Libp2pNetwork { .relay_invalidate_tokens_queries .insert(request_id, reply); } + SwarmCommand::RelayReadContent { + peer_id, + content_id, + version, + auth_token, + request_signature, + timestamp, + reply, + } => { + let request_id = swarm.behaviour_mut().request_response.send_request( + &peer_id, + ContentRequest::ReadContent { + content_id, + version, + auth_token, + request_signature, + timestamp, + }, + ); + pending.relay_read_queries.insert(request_id, reply); + } + SwarmCommand::RelayReadHistory { + peer_id, + content_id, + auth_token, + request_signature, + timestamp, + reply, + } => { + let request_id = swarm.behaviour_mut().request_response.send_request( + &peer_id, + ContentRequest::ReadHistory { + content_id, + auth_token, + request_signature, + timestamp, + }, + ); + pending.relay_history_queries.insert(request_id, reply); + } SwarmCommand::SendRelayResponse { channel, response } => { if let Err(e) = swarm .behaviour_mut() @@ -922,8 +1042,10 @@ impl Libp2pNetwork { domain_event.event_type() ); + // Bind to the authenticated publisher, not the peer + // that forwarded it to us — see `ReceivedEvent::source`. let received = ReceivedEvent { - source: propagation_source.to_string(), + source: message.source.map(|p| p.to_string()), event: domain_event, }; @@ -1013,6 +1135,12 @@ impl Libp2pNetwork { if let Some(reply) = pending.relay_invalidate_tokens_queries.remove(&request_id) { let _ = reply.send(Err(anyhow::anyhow!("{}", err_msg))); } + if let Some(reply) = pending.relay_read_queries.remove(&request_id) { + let _ = reply.send(Err(RelayReadError::other(&err_msg))); + } + if let Some(reply) = pending.relay_history_queries.remove(&request_id) { + let _ = reply.send(Err(RelayReadError::other(&err_msg))); + } } _ => {} } @@ -1056,8 +1184,11 @@ impl Libp2pNetwork { // binding `sender_peer == bs.creator_node_id` stops creator // impersonation. Requiring that `local_peer` appear in // `member_nodes` stops a malicious peer from fabricating a - // network on an unrelated victim node. The residual risk - // matches the Gossipsub ContentCreated trust model. + // network on an unrelated victim node. This mirrors the + // Gossipsub `ContentCreated` path, which binds the + // authenticated publisher to the creator it claims; the + // residual risk is the same — an authenticated creator can + // still declare a member set of its choosing. if bs.creator_node_id != sender_peer { return Err(format!( "bootstrap creator_node_id {} does not match sender {}", @@ -1148,7 +1279,7 @@ impl Libp2pNetwork { }; let response = if channels.relay_tx.send(relay_req).await.is_ok() { match reply_rx.await { - Ok(Ok(())) => ContentResponse::UpdateResult { + Ok(Ok(_)) => ContentResponse::UpdateResult { content_id, success: true, }, @@ -1195,7 +1326,7 @@ impl Libp2pNetwork { }; let response = if channels.relay_tx.send(relay_req).await.is_ok() { match reply_rx.await { - Ok(Ok(())) => ContentResponse::DeleteResult { + Ok(Ok(_)) => ContentResponse::DeleteResult { content_id, success: true, }, @@ -1242,7 +1373,7 @@ impl Libp2pNetwork { }; let response = if channels.relay_tx.send(relay_req).await.is_ok() { match reply_rx.await { - Ok(Ok(())) => ContentResponse::InvalidateTokensResult { + Ok(Ok(_)) => ContentResponse::InvalidateTokensResult { content_id, success: true, }, @@ -1265,6 +1396,117 @@ impl Libp2pNetwork { }); return; } + ContentRequest::ReadContent { + content_id, + version, + auth_token, + request_signature, + timestamp, + } => { + info!( + "Received relayed ReadContent for {} from {}", + content_id, peer + ); + let channels = relay_channels.clone(); + tokio::spawn(async move { + let (reply_tx, reply_rx) = oneshot::channel(); + let relay_req = RelayRequest { + kind: RelayRequestKind::ReadContent { + content_id: content_id.clone(), + version, + auth_token, + request_signature, + timestamp, + }, + reply: reply_tx, + }; + let response = if channels.relay_tx.send(relay_req).await.is_ok() { + match reply_rx.await { + Ok(Ok(RelayOutcome::Data { data, version })) => { + ContentResponse::ContentData { + content_id, + data, + version, + } + } + Ok(Ok(_)) => ContentResponse::Error { + message: "Unexpected relay outcome for ReadContent".to_string(), + }, + Ok(Err(e)) => ContentResponse::ReadFailed { + content_id, + kind: relay_read_error_kind(&e), + message: e.to_string(), + }, + Err(_) => ContentResponse::Error { + message: "Relay handler dropped".to_string(), + }, + } + } else { + ContentResponse::Error { + message: "Relay channel closed".to_string(), + } + }; + let _ = channels + .command_tx + .send(SwarmCommand::SendRelayResponse { channel, response }) + .await; + }); + return; + } + ContentRequest::ReadHistory { + content_id, + auth_token, + request_signature, + timestamp, + } => { + info!( + "Received relayed ReadHistory for {} from {}", + content_id, peer + ); + let channels = relay_channels.clone(); + tokio::spawn(async move { + let (reply_tx, reply_rx) = oneshot::channel(); + let relay_req = RelayRequest { + kind: RelayRequestKind::ReadHistory { + content_id: content_id.clone(), + auth_token, + request_signature, + timestamp, + }, + reply: reply_tx, + }; + let response = if channels.relay_tx.send(relay_req).await.is_ok() { + match reply_rx.await { + Ok(Ok(RelayOutcome::History { versions })) => { + ContentResponse::HistoryData { + content_id, + versions, + } + } + Ok(Ok(_)) => ContentResponse::Error { + message: "Unexpected relay outcome for ReadHistory".to_string(), + }, + Ok(Err(e)) => ContentResponse::ReadFailed { + content_id, + kind: relay_read_error_kind(&e), + message: e.to_string(), + }, + Err(_) => ContentResponse::Error { + message: "Relay handler dropped".to_string(), + }, + } + } else { + ContentResponse::Error { + message: "Relay channel closed".to_string(), + } + }; + let _ = channels + .command_tx + .send(SwarmCommand::SendRelayResponse { channel, response }) + .await; + }); + return; + } _ => {} } @@ -1447,7 +1689,9 @@ impl Libp2pNetwork { // Relay variants already handled above and returned early ContentRequest::UpdateContent { .. } | ContentRequest::DeleteContent { .. } - | ContentRequest::InvalidateTokens { .. } => unreachable!(), + | ContentRequest::InvalidateTokens { .. } + | ContentRequest::ReadContent { .. } + | ContentRequest::ReadHistory { .. } => unreachable!(), }; if let Err(e) = swarm @@ -1593,6 +1837,62 @@ impl Libp2pNetwork { let _ = reply.send(Err(anyhow::anyhow!("Unexpected response type"))); } } + return; + } + + // Handle relay read (data) response + if let Some(reply) = pending.relay_read_queries.remove(&request_id) { + match response { + ContentResponse::ContentData { data, version, .. } => { + let _ = reply.send(Ok((data, version))); + } + ContentResponse::ReadFailed { kind, message, .. } => { + let _ = reply.send(Err(RelayReadError { kind, message })); + } + ContentResponse::NotFound { content_id } => { + let _ = reply.send(Err(RelayReadError { + kind: RelayReadErrorKind::NotFound, + message: format!("Content not found: {}", content_id), + })); + } + ContentResponse::Error { message } => { + let _ = reply.send(Err(RelayReadError::other(format!( + "Relay read error: {}", + message + )))); + } + _ => { + let _ = reply.send(Err(RelayReadError::other("Unexpected response type"))); + } + } + return; + } + + // Handle relay read (history) response + if let Some(reply) = pending.relay_history_queries.remove(&request_id) { + match response { + ContentResponse::HistoryData { versions, .. } => { + let _ = reply.send(Ok(versions)); + } + ContentResponse::ReadFailed { kind, message, .. } => { + let _ = reply.send(Err(RelayReadError { kind, message })); + } + ContentResponse::NotFound { content_id } => { + let _ = reply.send(Err(RelayReadError { + kind: RelayReadErrorKind::NotFound, + message: format!("Content not found: {}", content_id), + })); + } + ContentResponse::Error { message } => { + let _ = reply.send(Err(RelayReadError::other(format!( + "Relay read error: {}", + message + )))); + } + _ => { + let _ = reply.send(Err(RelayReadError::other("Unexpected response type"))); + } + } } } @@ -2138,6 +2438,68 @@ impl PeerNetwork for Libp2pNetwork { .map_err(|_| anyhow::anyhow!("Failed to receive response"))? } + async fn relay_read_content( + &self, + peer_id: &str, + content_id: &str, + version: Option<&str>, + auth_token: &str, + request_signature: &[u8], + timestamp: Option, + ) -> std::result::Result<(Vec, String), RelayReadError> { + let peer_id = PeerId::from_str(peer_id) + .map_err(|_| RelayReadError::other(format!("Invalid peer ID: {}", peer_id)))?; + + let (tx, rx) = oneshot::channel(); + self.command_tx + .send(SwarmCommand::RelayReadContent { + peer_id, + content_id: content_id.to_string(), + version: version.map(ToOwned::to_owned), + auth_token: auth_token.to_string(), + request_signature: request_signature.to_vec(), + timestamp, + reply: tx, + }) + .await + .map_err(|_| RelayReadError::other("Failed to send command"))?; + + tokio::time::timeout(PEER_NETWORK_TIMEOUT, rx) + .await + .map_err(|_| RelayReadError::other("relay_read_content timed out"))? + .map_err(|_| RelayReadError::other("Failed to receive response"))? + } + + async fn relay_read_history( + &self, + peer_id: &str, + content_id: &str, + auth_token: &str, + request_signature: &[u8], + timestamp: Option, + ) -> std::result::Result, RelayReadError> { + let peer_id = PeerId::from_str(peer_id) + .map_err(|_| RelayReadError::other(format!("Invalid peer ID: {}", peer_id)))?; + + let (tx, rx) = oneshot::channel(); + self.command_tx + .send(SwarmCommand::RelayReadHistory { + peer_id, + content_id: content_id.to_string(), + auth_token: auth_token.to_string(), + request_signature: request_signature.to_vec(), + timestamp, + reply: tx, + }) + .await + .map_err(|_| RelayReadError::other("Failed to send command"))?; + + tokio::time::timeout(PEER_NETWORK_TIMEOUT, rx) + .await + .map_err(|_| RelayReadError::other("relay_read_history timed out"))? + .map_err(|_| RelayReadError::other("Failed to receive response"))? + } + async fn connected_peer_count(&self) -> usize { self.connected_peers.read().await.len() } diff --git a/monas-state-node/src/infrastructure/network/protocol.rs b/monas-state-node/src/infrastructure/network/protocol.rs index 63b98ae..ef9ac5b 100644 --- a/monas-state-node/src/infrastructure/network/protocol.rs +++ b/monas-state-node/src/infrastructure/network/protocol.rs @@ -6,7 +6,7 @@ use serde::{Deserialize, Serialize}; -pub use crate::port::peer_network::PushBootstrap; +pub use crate::port::peer_network::{PushBootstrap, RelayReadErrorKind}; /// Protocol name for capacity queries. pub const CAPACITY_PROTOCOL: &str = "/monas/capacity/1.0.0"; @@ -66,6 +66,30 @@ pub enum ContentRequest { request_signature: Vec, timestamp: Option, }, + /// Relay a content-data read to a member node. + /// + /// Sent by a node that does not replicate the content (e.g. the + /// gateway-facing node) so a member can serve the read. The member + /// re-authenticates the original caller from `auth_token` / + /// `request_signature` and enforces the content's access policy before + /// returning any data. + ReadContent { + content_id: String, + /// `None` reads the latest version; `Some(v)` a specific version CID. + version: Option, + auth_token: String, + request_signature: Vec, + timestamp: Option, + }, + /// Relay a version-history read to a member node. + /// + /// Same authentication contract as [`ContentRequest::ReadContent`]. + ReadHistory { + content_id: String, + auth_token: String, + request_signature: Vec, + timestamp: Option, + }, } /// Response types for the content protocol. @@ -101,6 +125,19 @@ pub enum ContentResponse { DeleteResult { content_id: String, success: bool }, /// Response to relayed invalidate_tokens request. InvalidateTokensResult { content_id: String, success: bool }, + /// Response to a relayed history read. + HistoryData { + content_id: String, + versions: Vec, + }, + /// Failure of a relayed read, carrying the member's typed verdict so the + /// relaying node can map it back to 401/403/404 without parsing message + /// text. + ReadFailed { + content_id: String, + kind: RelayReadErrorKind, + message: String, + }, /// Content not found. NotFound { content_id: String }, /// Error response. diff --git a/monas-state-node/src/infrastructure/persistence/sled_public_key_repository.rs b/monas-state-node/src/infrastructure/persistence/sled_public_key_repository.rs index 09065b1..785095f 100644 --- a/monas-state-node/src/infrastructure/persistence/sled_public_key_repository.rs +++ b/monas-state-node/src/infrastructure/persistence/sled_public_key_repository.rs @@ -302,9 +302,20 @@ mod tests { repo.flush().await.unwrap(); } - // Open new repository instance and verify key persists + // Open new repository instance and verify key persists. + // sled releases its file lock asynchronously on Drop, so an immediate + // reopen can transiently fail with WouldBlock on slow CI runners — + // retry briefly instead of failing the test on that race. { - let repo = SledPublicKeyRepository::open(temp_dir.path()).unwrap(); + let mut repo = SledPublicKeyRepository::open(temp_dir.path()); + for _ in 0..20 { + if repo.is_ok() { + break; + } + tokio::time::sleep(std::time::Duration::from_millis(50)).await; + repo = SledPublicKeyRepository::open(temp_dir.path()); + } + let repo = repo.unwrap(); let retrieved = repo.get_public_key(&node_id).await.unwrap(); assert!(retrieved.is_some()); assert_eq!(retrieved.unwrap(), public_key); diff --git a/monas-state-node/src/port/content_repository.rs b/monas-state-node/src/port/content_repository.rs index fe0a4af..dcbf0e4 100644 --- a/monas-state-node/src/port/content_repository.rs +++ b/monas-state-node/src/port/content_repository.rs @@ -109,14 +109,19 @@ pub trait ContentRepository: Send + Sync { async fn get_latest_with_version(&self, genesis_cid: &str) -> Result, String)>>; - /// Get content at a specific version. + /// Get content at a specific version, scoped to one content series. /// /// # Arguments + /// * `genesis_cid` - The genesis CID of the content the caller is + /// authorized for. The version MUST belong to this series — a raw CID + /// lookup would let a caller authorized for one content read any other + /// content's versions. /// * `version_cid` - The specific version CID /// /// # Returns - /// The content data at that version, or None if not found. - async fn get_version(&self, version_cid: &str) -> Result>>; + /// The content data at that version, or None if the version does not + /// exist or belongs to a different content series. + async fn get_version(&self, genesis_cid: &str, version_cid: &str) -> Result>>; /// Get the version history of content. /// diff --git a/monas-state-node/src/port/peer_network.rs b/monas-state-node/src/port/peer_network.rs index 7514d96..5686443 100644 --- a/monas-state-node/src/port/peer_network.rs +++ b/monas-state-node/src/port/peer_network.rs @@ -28,6 +28,43 @@ pub struct PushBootstrap { pub created_at: u64, } +/// Machine-readable failure category of a relayed read, decided by the +/// responding member. +/// +/// Crosses the wire inside `ContentResponse::ReadFailed` so the relaying +/// node can return the member's verdict (401/403/404) as a typed error +/// instead of re-deriving it from error-message text. +#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)] +pub enum RelayReadErrorKind { + /// The member does not hold the content/version (404). + NotFound, + /// The member could not authenticate the forwarded caller (401). + AuthenticationFailed, + /// The member authenticated the caller but denied read access (403). + AuthorizationFailed, + /// Transport failures, timeouts, or unclassified member errors. Unlike + /// the verdict kinds above, this does not represent a member decision. + Other, +} + +/// Failure of a relayed read: the member's verdict (`kind`) plus its +/// human-readable message. +#[derive(Debug, Clone, thiserror::Error)] +#[error("{message}")] +pub struct RelayReadError { + pub kind: RelayReadErrorKind, + pub message: String, +} + +impl RelayReadError { + pub fn other(message: impl Into) -> Self { + Self { + kind: RelayReadErrorKind::Other, + message: message.into(), + } + } +} + /// Abstract interface for peer-to-peer network operations. /// /// This trait provides methods for: @@ -171,6 +208,34 @@ pub trait PeerNetwork: Send + Sync { timestamp: Option, ) -> Result; + /// Relay a content-data read to a member node. + /// + /// Used when a node that does not replicate the content receives a read + /// request. The member re-authenticates the original caller (token + + /// request signature) and enforces the access policy before serving. + /// `version: None` reads the latest version. Returns `(data, version)`. + async fn relay_read_content( + &self, + peer_id: &str, + content_id: &str, + version: Option<&str>, + auth_token: &str, + request_signature: &[u8], + timestamp: Option, + ) -> std::result::Result<(Vec, String), RelayReadError>; + + /// Relay a version-history read to a member node. + /// + /// Same authentication contract as [`PeerNetwork::relay_read_content`]. + async fn relay_read_history( + &self, + peer_id: &str, + content_id: &str, + auth_token: &str, + request_signature: &[u8], + timestamp: Option, + ) -> std::result::Result, RelayReadError>; + // ========== Monitoring Methods ========== /// Get the number of currently connected peers. diff --git a/monas-state-node/src/presentation/http_api.rs b/monas-state-node/src/presentation/http_api.rs index 2dc4129..889fe16 100644 --- a/monas-state-node/src/presentation/http_api.rs +++ b/monas-state-node/src/presentation/http_api.rs @@ -613,66 +613,84 @@ async fn verify_read_access( let request_sig = extract_request_signature(headers); let timestamp = extract_request_timestamp(headers); - // Authenticate the caller - let identity = state - .authenticate_for_read(&token, request_sig.as_deref(), timestamp) + state + .authorize_read(&token, request_sig.as_deref(), timestamp, content_id) .await - .map_err(|e| { - ( - StatusCode::UNAUTHORIZED, - Json(ErrorResponse { - error: format!("Authentication failed: {}", e), - }), - ) - .into_response() - })?; - - // Check access policy for read permission - let crdt_repo = state.crdt_repo(); - if let Ok(Some(policy)) = crdt_repo.get_access_policy(content_id).await { - // Owner always has access - if policy.is_owner(&identity) { - return Ok(()); - } + .map_err(|e| e.into_response()) +} - // Non-owner: needs AuthToken-based authorization - // The Bearer token is used as-is for JWT-based auth - let request_signature = extract_request_signature(headers); - let authz_request = crate::port::authorization_service::AuthorizationRequest { - identity, - resource: crate::domain::value_objects::ContentId::new(content_id.to_string()) - .map_err(|_| { - ( - StatusCode::BAD_REQUEST, - Json(ErrorResponse { - error: "Invalid content ID".to_string(), - }), - ) - .into_response() - })?, - capability: crate::domain::auth_capability::AuthCapability::ReadContent, - token: Some(token), - request_signature, - }; +/// Serve a read for content this node does not replicate by relaying it — +/// auth material included — to a member node (bug #93 hardened read path). +async fn relay_read_data_response( + state: &AppState, + headers: &HeaderMap, + content_id: &str, + version: Option<&str>, +) -> Response { + let Some(token) = extract_auth_token(headers) else { + return ( + StatusCode::UNAUTHORIZED, + Json(ErrorResponse { + error: "Authorization header is required".to_string(), + }), + ) + .into_response(); + }; + let request_sig = extract_request_signature(headers); + let timestamp = extract_request_timestamp(headers); - if let Some(authz_service) = state.authz_service() { - match authz_service.authorize(&authz_request).await { - Ok(result) if result.is_granted() => return Ok(()), - _ => {} - } + match state + .relay_read_data( + content_id, + version, + &token, + request_sig.as_deref(), + timestamp, + ) + .await + { + Ok((data, served_version)) => { + let encoded = base64::engine::general_purpose::STANDARD.encode(&data); + Json(ContentDataResponse { + content_id: content_id.to_string(), + data: encoded, + version: Some(served_version), + }) + .into_response() } + Err(e) => e.into_response(), + } +} - return Err(( - StatusCode::FORBIDDEN, +/// History counterpart of [`relay_read_data_response`]. +async fn relay_read_history_response( + state: &AppState, + headers: &HeaderMap, + content_id: &str, +) -> Response { + let Some(token) = extract_auth_token(headers) else { + return ( + StatusCode::UNAUTHORIZED, Json(ErrorResponse { - error: "Insufficient permissions: read access required".to_string(), + error: "Authorization header is required".to_string(), }), ) - .into_response()); - } - // If no policy exists, allow access (content may not have a policy yet) + .into_response(); + }; + let request_sig = extract_request_signature(headers); + let timestamp = extract_request_timestamp(headers); - Ok(()) + match state + .relay_read_history(content_id, &token, request_sig.as_deref(), timestamp) + .await + { + Ok(versions) => Json(ContentHistoryResponse { + content_id: content_id.to_string(), + versions, + }) + .into_response(), + Err(e) => e.into_response(), + } } /// Get content data from CRDT repository. @@ -684,11 +702,14 @@ async fn get_content_data( headers: HeaderMap, Query(query): Query, ) -> impl IntoResponse { - // Bug #93: if this node holds no local state for the content (it is neither - // the creator nor a member), pull it from a member first so the read below - // and the access-policy check both see the real data. Best-effort: on - // failure we fall through to the normal local read (which 404s as before). - let _ = state.ensure_content_local(&content_id).await; + // Bug #93: this node may not replicate the content (it is neither the + // creator nor a member). Relay the read — auth material included — to a + // member node, which re-authenticates the caller against the real access + // policy and serves the data. Operations are never pulled to this node. + if !state.has_local_content(&content_id).await { + return relay_read_data_response(&state, &headers, &content_id, query.version.as_deref()) + .await; + } if let Err(response) = verify_read_access(&state, &headers, &content_id).await { return response; @@ -698,7 +719,7 @@ async fn get_content_data( // Get data based on version parameter let data_result = if let Some(version) = &query.version { - crdt_repo.get_version(version).await + crdt_repo.get_version(&content_id, version).await } else { crdt_repo.get_latest(&content_id).await }; @@ -741,8 +762,10 @@ async fn get_content_history( Path(content_id): Path, headers: HeaderMap, ) -> impl IntoResponse { - // Bug #93: pull content from a member if we hold none locally (best-effort). - let _ = state.ensure_content_local(&content_id).await; + // Bug #93: relay the read to a member when we don't replicate the content. + if !state.has_local_content(&content_id).await { + return relay_read_history_response(&state, &headers, &content_id).await; + } if let Err(response) = verify_read_access(&state, &headers, &content_id).await { return response; @@ -777,8 +800,10 @@ async fn get_content_version( Path((content_id, version)): Path<(String, String)>, headers: HeaderMap, ) -> impl IntoResponse { - // Bug #93: pull content from a member if we hold none locally (best-effort). - let _ = state.ensure_content_local(&content_id).await; + // Bug #93: relay the read to a member when we don't replicate the content. + if !state.has_local_content(&content_id).await { + return relay_read_data_response(&state, &headers, &content_id, Some(&version)).await; + } if let Err(response) = verify_read_access(&state, &headers, &content_id).await { return response; @@ -786,7 +811,7 @@ async fn get_content_version( let crdt_repo = state.crdt_repo(); - match crdt_repo.get_version(&version).await { + match crdt_repo.get_version(&content_id, &version).await { Ok(Some(data)) => { let encoded = base64::engine::general_purpose::STANDARD.encode(&data); Json(ContentDataResponse { diff --git a/monas-state-node/src/test_utils.rs b/monas-state-node/src/test_utils.rs index 07abec5..5433bf0 100644 --- a/monas-state-node/src/test_utils.rs +++ b/monas-state-node/src/test_utils.rs @@ -9,7 +9,7 @@ use crate::domain::events::Event; use crate::domain::state_node::NodeSnapshot; use crate::port::content_repository::{CommitResult, ContentRepository, SerializedOperation}; use crate::port::event_publisher::EventPublisher; -use crate::port::peer_network::PeerNetwork; +use crate::port::peer_network::{PeerNetwork, RelayReadError, RelayReadErrorKind}; use crate::port::persistence::{PersistentContentRepository, PersistentNodeRegistry}; use anyhow::Result; use async_trait::async_trait; @@ -24,6 +24,9 @@ use tokio::sync::Mutex; /// Type alias for published events storage. pub type PublishedEvents = Arc)>>>; +/// Configurable result of a mocked relayed data read: `(data, version)`. +pub type MockRelayReadData = Arc, String)>>>; + /// Mock implementation of PeerNetwork for testing. #[derive(Default)] pub struct MockPeerNetwork { @@ -47,6 +50,18 @@ pub struct MockPeerNetwork { pub relay_update_peers: Arc>>, pub relay_delete_peers: Arc>>, pub relay_invalidate_tokens_peers: Arc>>, + /// When `Some`, relayed data reads succeed with this payload; when `None` + /// they fail with a typed NotFound verdict (the common test default). + pub relay_read_data_result: MockRelayReadData, + /// Same, for relayed history reads. + pub relay_read_history_result: Arc>>>, + /// When `Some`, relayed data reads fail with exactly this typed error + /// (takes precedence over `relay_read_data_result`). + pub relay_read_data_error: Arc>>, + /// `since_version` argument of each `fetch_operations` call, in order. + /// Lets tests assert whether a sync requested the full history + /// (`None`) or an incremental fetch (`Some(version)`). + pub fetch_operations_since: Arc>>>, } impl MockPeerNetwork { @@ -68,6 +83,31 @@ impl MockPeerNetwork { relay_update_peers: Arc::new(Mutex::new(Vec::new())), relay_delete_peers: Arc::new(Mutex::new(Vec::new())), relay_invalidate_tokens_peers: Arc::new(Mutex::new(Vec::new())), + relay_read_data_result: Arc::new(Mutex::new(None)), + relay_read_history_result: Arc::new(Mutex::new(None)), + relay_read_data_error: Arc::new(Mutex::new(None)), + fetch_operations_since: Arc::new(Mutex::new(Vec::new())), + } + } + + pub fn with_relay_read_error(self, error: RelayReadError) -> Self { + Self { + relay_read_data_error: Arc::new(Mutex::new(Some(error))), + ..self + } + } + + pub fn with_relay_read_data(self, data: Vec, version: &str) -> Self { + Self { + relay_read_data_result: Arc::new(Mutex::new(Some((data, version.to_string())))), + ..self + } + } + + pub fn with_relay_read_history(self, versions: Vec) -> Self { + Self { + relay_read_history_result: Arc::new(Mutex::new(Some(versions))), + ..self } } @@ -176,8 +216,12 @@ impl PeerNetwork for MockPeerNetwork { &self, _peer_id: &str, _genesis_cid: &str, - _since_version: Option<&str>, + since_version: Option<&str>, ) -> Result> { + self.fetch_operations_since + .lock() + .await + .push(since_version.map(ToOwned::to_owned)); Ok(self.fetched_operations.lock().await.clone()) } @@ -265,6 +309,44 @@ impl PeerNetwork for MockPeerNetwork { .unwrap_or(true)) } + async fn relay_read_content( + &self, + _peer_id: &str, + content_id: &str, + _version: Option<&str>, + _auth_token: &str, + _request_signature: &[u8], + _timestamp: Option, + ) -> std::result::Result<(Vec, String), RelayReadError> { + if let Some(error) = self.relay_read_data_error.lock().await.clone() { + return Err(error); + } + match self.relay_read_data_result.lock().await.clone() { + Some(result) => Ok(result), + None => Err(RelayReadError { + kind: RelayReadErrorKind::NotFound, + message: format!("Content not found: {}", content_id), + }), + } + } + + async fn relay_read_history( + &self, + _peer_id: &str, + content_id: &str, + _auth_token: &str, + _request_signature: &[u8], + _timestamp: Option, + ) -> std::result::Result, RelayReadError> { + match self.relay_read_history_result.lock().await.clone() { + Some(versions) => Ok(versions), + None => Err(RelayReadError { + kind: RelayReadErrorKind::NotFound, + message: format!("Content not found: {}", content_id), + }), + } + } + async fn connected_peer_count(&self) -> usize { 0 } @@ -322,6 +404,9 @@ pub struct MockContentRepository { pub operations: Arc>>, pub next_cid: Arc>, pub access_policies: Arc>>, + /// When true, `get_access_policy` fails — used to test that read + /// authorization fails closed on policy-store errors. + pub access_policy_error: Arc>, } impl MockContentRepository { @@ -332,6 +417,7 @@ impl MockContentRepository { operations: Arc::new(Mutex::new(Vec::new())), next_cid: Arc::new(Mutex::new(1)), access_policies: Arc::new(Mutex::new(HashMap::new())), + access_policy_error: Arc::new(Mutex::new(false)), } } } @@ -415,14 +501,12 @@ impl ContentRepository for MockContentRepository { } } - async fn get_version(&self, version_cid: &str) -> Result>> { - // For simplicity, return the first content that matches + async fn get_version(&self, genesis_cid: &str, version_cid: &str) -> Result>> { + // Scoped to the given series, mirroring the real implementation. let contents = self.contents.lock().await; - for genesis_cid in contents.keys() { - if let Some(history) = self.history.lock().await.get(genesis_cid) { - if history.contains(&version_cid.to_string()) { - return Ok(contents.get(genesis_cid).cloned()); - } + if let Some(history) = self.history.lock().await.get(genesis_cid) { + if history.contains(&version_cid.to_string()) { + return Ok(contents.get(genesis_cid).cloned()); } } Ok(None) @@ -467,6 +551,9 @@ impl ContentRepository for MockContentRepository { } async fn get_access_policy(&self, genesis_cid: &str) -> Result> { + if *self.access_policy_error.lock().await { + return Err(anyhow::anyhow!("policy store unavailable")); + } Ok(self.access_policies.lock().await.get(genesis_cid).cloned()) }