From 991a63ecf20d26a0086114a81117c8ac46fb631a Mon Sep 17 00:00:00 2001 From: misrasaurabh1 Date: Wed, 22 Jul 2026 11:36:37 -0700 Subject: [PATCH 1/3] Handle Notion database roots returned as type mismatches --- crates/locality-notion/src/client.rs | 159 ++++++++++++++++++++++++--- 1 file changed, 144 insertions(+), 15 deletions(-) diff --git a/crates/locality-notion/src/client.rs b/crates/locality-notion/src/client.rs index 320169c1..b7905da0 100644 --- a/crates/locality-notion/src/client.rs +++ b/crates/locality-notion/src/client.rs @@ -329,6 +329,13 @@ enum NotionRetryClass { Mutation, } +#[derive(Clone, Copy, Debug, Default, PartialEq, Eq)] +enum NotionResponseInterpretation { + #[default] + Default, + PageLookup, +} + impl HttpNotionApi { pub fn new(config: NotionConfig) -> Self { let client = notion_http_client_builder() @@ -339,6 +346,25 @@ impl HttpNotionApi { } fn get_json(&self, path: &str, query: &[(&str, String)]) -> LocalityResult + where + T: DeserializeOwned, + { + self.get_json_with_interpretation(path, query, NotionResponseInterpretation::Default) + } + + fn get_page_json(&self, path: &str) -> LocalityResult + where + T: DeserializeOwned, + { + self.get_json_with_interpretation(path, &[], NotionResponseInterpretation::PageLookup) + } + + fn get_json_with_interpretation( + &self, + path: &str, + query: &[(&str, String)], + response_interpretation: NotionResponseInterpretation, + ) -> LocalityResult where T: DeserializeOwned, { @@ -353,18 +379,24 @@ impl HttpNotionApi { .map(|(key, value)| ((*key).to_string(), value.clone())) .collect::>(); - self.send_request_with_retry("GET", path, NotionRetryClass::ReadSafe, || { - let mut request = self - .client - .get(&url) - .bearer_auth(&token) - .header("Notion-Version", DEFAULT_NOTION_VERSION); - - for (key, value) in &query { - request = request.query(&[(key.as_str(), value.as_str())]); - } - request - }) + self.send_request_with_retry_and_interpretation( + "GET", + path, + NotionRetryClass::ReadSafe, + response_interpretation, + || { + let mut request = self + .client + .get(&url) + .bearer_auth(&token) + .header("Notion-Version", DEFAULT_NOTION_VERSION); + + for (key, value) in &query { + request = request.query(&[(key.as_str(), value.as_str())]); + } + request + }, + ) } fn post_json(&self, path: &str, body: impl Serialize) -> LocalityResult @@ -532,6 +564,26 @@ impl HttpNotionApi { method: &str, path: &str, retry_class: NotionRetryClass, + build_request: impl FnMut() -> reqwest::blocking::RequestBuilder, + ) -> LocalityResult + where + T: DeserializeOwned, + { + self.send_request_with_retry_and_interpretation( + method, + path, + retry_class, + NotionResponseInterpretation::Default, + build_request, + ) + } + + fn send_request_with_retry_and_interpretation( + &self, + method: &str, + path: &str, + retry_class: NotionRetryClass, + response_interpretation: NotionResponseInterpretation, mut build_request: impl FnMut() -> reqwest::blocking::RequestBuilder, ) -> LocalityResult where @@ -575,6 +627,15 @@ impl HttpNotionApi { let body = response .text() .unwrap_or_else(|error| format!("")); + if response_interpretation == NotionResponseInterpretation::PageLookup + && notion_page_lookup_reports_database(status, &body) + { + // Notion reports an object-kind mismatch as HTTP 400 rather + // than 404. Present it as a page miss only at this exact + // boundary so explicit-root traversal can try the database + // endpoint without weakening other validation failures. + return Err(LocalityError::RemoteNotFound(body)); + } if is_retryable_notion_http_status(status, retry_class) && attempt < DEFAULT_NOTION_RATE_LIMIT_RETRIES { @@ -673,6 +734,20 @@ fn is_retryable_notion_transport_error(error: &reqwest::Error) -> bool { error.is_timeout() || error.is_connect() || error.is_request() || error.is_body() } +fn notion_page_lookup_reports_database(status: StatusCode, body: &str) -> bool { + if status != StatusCode::BAD_REQUEST { + return false; + } + let Ok(error) = serde_json::from_str::(body) else { + return false; + }; + error.get("code").and_then(Value::as_str) == Some("validation_error") + && error + .get("message") + .and_then(Value::as_str) + .is_some_and(|message| message.contains(" is a database, not a page")) +} + fn retry_after_header(headers: &HeaderMap) -> Option { headers .get(reqwest::header::RETRY_AFTER)? @@ -710,7 +785,7 @@ impl NotionApi for HttpNotionApi { } fn retrieve_page(&self, page_id: &str) -> LocalityResult { - self.get_json(&format!("/v1/pages/{page_id}"), &[]) + self.get_page_json(&format!("/v1/pages/{page_id}")) } fn retrieve_database(&self, database_id: &str) -> LocalityResult { @@ -879,8 +954,9 @@ fn data_source_search_body(start_cursor: Option<&str>) -> Value { #[cfg(test)] mod tests { use super::{ - HttpNotionApi, NotionRetryClass, data_source_search_body, notion_http_client_builder, - notion_network_config, rate_limit_backoff, retry_after_header, + HttpNotionApi, NotionResponseInterpretation, NotionRetryClass, data_source_search_body, + notion_http_client_builder, notion_network_config, notion_page_lookup_reports_database, + rate_limit_backoff, retry_after_header, }; use locality_core::LocalityError; use reqwest::header::{HeaderMap, HeaderValue, RETRY_AFTER}; @@ -931,6 +1007,59 @@ mod tests { assert_eq!(config.retry.max_backoff, Duration::from_secs(16)); } + #[test] + fn page_lookup_maps_exact_database_kind_mismatch_to_page_miss() { + let listener = TcpListener::bind("127.0.0.1:0").expect("bind local server"); + let url = format!("http://{}/database-root", listener.local_addr().unwrap()); + let body = r#"{"object":"error","status":400,"code":"validation_error","message":"Provided ID 4614fba4-9bdf-45e0-a006-4f91dca082f1 is a database, not a page. Use the retrieve database API instead.","request_id":"request-1"}"#; + let response_body = body.as_bytes().to_vec(); + let server = thread::spawn(move || { + let (mut stream, _) = listener.accept().expect("accept request"); + read_http_request_headers(&mut stream); + write!( + stream, + "HTTP/1.1 400 Bad Request\r\nContent-Type: application/json\r\nContent-Length: {}\r\nConnection: close\r\n\r\n", + response_body.len() + ) + .expect("write headers"); + stream.write_all(&response_body).expect("write body"); + }); + let api = HttpNotionApi { + config: crate::NotionConfig::default(), + client: notion_http_client_builder() + .timeout(Duration::from_millis(500)) + .build() + .expect("build client"), + }; + + let error = api + .send_request_with_retry_and_interpretation::( + "GET", + "/v1/pages/database-root", + NotionRetryClass::ReadSafe, + NotionResponseInterpretation::PageLookup, + || api.client.get(&url), + ) + .expect_err("database kind mismatch is a page miss"); + + server.join().expect("join server"); + assert_eq!(error, LocalityError::RemoteNotFound(body.to_string())); + } + + #[test] + fn page_lookup_does_not_reclassify_other_bad_requests() { + let other_validation = + r#"{"code":"validation_error","message":"Provided page ID is invalid."}"#; + assert!(!notion_page_lookup_reports_database( + reqwest::StatusCode::BAD_REQUEST, + other_validation + )); + assert!(!notion_page_lookup_reports_database( + reqwest::StatusCode::NOT_FOUND, + r#"{"code":"validation_error","message":"ID is a database, not a page"}"# + )); + } + #[test] fn send_request_retries_transient_timeout_before_returning_response() { let listener = TcpListener::bind("127.0.0.1:0").expect("bind local server"); From ca5226d16869e5c7ed7d3115f99e89f5a97a1df8 Mon Sep 17 00:00:00 2001 From: misrasaurabh1 Date: Wed, 22 Jul 2026 14:29:30 -0700 Subject: [PATCH 2/3] Aggregate paginated bootstrap candidates --- .../src/synchronize_project.rs | 295 ++++++++++- .../tests/synchronize_project.rs | 492 +++++++++++++++++- 2 files changed, 784 insertions(+), 3 deletions(-) diff --git a/crates/locality-engine/src/synchronize_project.rs b/crates/locality-engine/src/synchronize_project.rs index e7f7cd69..f1248e3c 100644 --- a/crates/locality-engine/src/synchronize_project.rs +++ b/crates/locality-engine/src/synchronize_project.rs @@ -9,8 +9,8 @@ use std::collections::{BTreeMap, BTreeSet}; use locality_connector::{ Connector, PortableArtifactKey, PortableBootstrapRequest, PortableChangeBatch, PortableCompleteness, PortableContentArtifact, PortableEnumerateRequest, - PortableEnumerateResult, PortableFetchReason, PortableFetchRequest, PortableProjectionArtifact, - PortableRenderRequest, PortableSourceChange, PortableSyncRequest, + PortableEnumerateResult, PortableFetchReason, PortableFetchRequest, PortableIncompleteReason, + PortableProjectionArtifact, PortableRenderRequest, PortableSourceChange, PortableSyncRequest, portable_scope_root_remote_id, }; use locality_core::model::RemoteId; @@ -93,6 +93,32 @@ impl UnpersistedSynchronizationBatch { } } +/// Hard bounds for aggregating a paginated portable bootstrap. +/// +/// These limits apply to the aggregate rather than to an individual connector +/// request. `PortableBootstrapRequest::max_changes` remains the per-checkpoint +/// provider-work bound. +#[derive(Clone, Copy, Debug, PartialEq, Eq)] +pub struct BootstrapAggregationLimits { + pub max_checkpoints: usize, + pub max_total_changes: usize, + pub max_total_content_bytes: u64, +} + +impl BootstrapAggregationLimits { + fn validate(self) -> LocalityResult { + if self.max_checkpoints == 0 + || self.max_total_changes == 0 + || self.max_total_content_bytes == 0 + { + return Err(aggregation_error( + "portable bootstrap aggregation limits must be nonzero", + )); + } + Ok(self) + } +} + /// Run one bootstrap checkpoint through fetch and render. pub fn bootstrap_and_project( connector: &C, @@ -110,6 +136,271 @@ pub fn bootstrap_and_project( ) } +/// Run every checkpoint of one bounded bootstrap and return one deterministic +/// unpersisted candidate. +/// +/// `CheckpointContinuation` is pagination control flow, not a coverage gap. +/// Every other incomplete reason is retained and therefore continues to block +/// publication after the terminal checkpoint is reached. +pub fn bootstrap_and_project_to_completion( + connector: &C, + request: PortableBootstrapRequest, + format_version: u32, + limits: BootstrapAggregationLimits, +) -> LocalityResult { + let limits = limits.validate()?; + let expected_source_connection_id = request.source_connection_id.clone(); + let scope = request.scope; + let max_changes = request.max_changes; + let mut current_checkpoint = request.checkpoint; + let mut seen_checkpoints = BTreeSet::new(); + if let Some(checkpoint) = ¤t_checkpoint { + seen_checkpoints.insert(checkpoint_identity(checkpoint)); + } + let mut aggregate = BootstrapAggregate::new(expected_source_connection_id.clone()); + let mut checkpoint_count = 0_usize; + + loop { + if checkpoint_count >= limits.max_checkpoints { + return Err(aggregation_error( + "portable bootstrap aggregation exceeded its checkpoint limit", + )); + } + checkpoint_count += 1; + let page = bootstrap_and_project( + connector, + PortableBootstrapRequest { + source_connection_id: expected_source_connection_id.clone(), + scope: scope.clone(), + checkpoint: current_checkpoint.clone(), + max_changes, + }, + format_version, + ) + .map_err(|_| aggregation_error("portable bootstrap aggregation page failed"))?; + if page.source_connection_id != expected_source_connection_id { + return Err(aggregation_error( + "portable bootstrap aggregation changed source connection", + )); + } + + let (continuation, preserved_completeness) = + completeness_without_continuation(&page.completeness); + if continuation { + validate_continuation_checkpoint( + current_checkpoint.as_ref(), + &page.next_checkpoint, + &mut seen_checkpoints, + )?; + } + let next_checkpoint = page.next_checkpoint.clone(); + aggregate.push(page, preserved_completeness, limits)?; + + if !continuation { + return Ok(aggregate.finish(next_checkpoint)); + } + current_checkpoint = Some(next_checkpoint); + } +} + +struct BootstrapAggregate { + source_connection_id: SourceConnectionId, + observed_changes: BTreeMap, + source_versions: BTreeMap, + contents: BTreeMap, + projections: BTreeMap, + projection_paths: BTreeSet, + completeness: PortableCompleteness, + total_changes: usize, + total_content_bytes: u64, +} + +impl BootstrapAggregate { + fn new(source_connection_id: SourceConnectionId) -> Self { + Self { + source_connection_id, + observed_changes: BTreeMap::new(), + source_versions: BTreeMap::new(), + contents: BTreeMap::new(), + projections: BTreeMap::new(), + projection_paths: BTreeSet::new(), + completeness: PortableCompleteness::complete(), + total_changes: 0, + total_content_bytes: 0, + } + } + + fn push( + &mut self, + page: UnpersistedSynchronizationBatch, + preserved_completeness: PortableCompleteness, + limits: BootstrapAggregationLimits, + ) -> LocalityResult<()> { + for source in &page.source_versions { + if self + .source_versions + .contains_key(&source.source_object.remote_id) + { + return Err(aggregation_error( + "portable bootstrap aggregation repeated a source version", + )); + } + } + for change in &page.observed_changes { + if self + .observed_changes + .contains_key(&change.source_object.remote_id) + { + return Err(aggregation_error( + "portable bootstrap aggregation repeated an observed source", + )); + } + } + for projection in &page.projections { + if self.projections.contains_key(&projection.artifact_key) { + return Err(aggregation_error( + "portable bootstrap aggregation repeated a projection artifact", + )); + } + if self + .projection_paths + .contains(projection.logical_path.as_str()) + { + return Err(aggregation_error( + "portable bootstrap aggregation repeated a logical path", + )); + } + } + for content in &page.contents { + if self.contents.contains_key(&content.artifact_key) { + return Err(aggregation_error( + "portable bootstrap aggregation repeated a content artifact", + )); + } + } + + let total_changes = self + .total_changes + .checked_add(page.observed_changes.len()) + .ok_or_else(|| { + aggregation_error("portable bootstrap aggregation change count overflowed") + })?; + if total_changes > limits.max_total_changes { + return Err(aggregation_error( + "portable bootstrap aggregation exceeded its change limit", + )); + } + let page_content_bytes = page.contents.iter().try_fold(0_u64, |total, content| { + total.checked_add(content.byte_length).ok_or_else(|| { + aggregation_error("portable bootstrap aggregation content bytes overflowed") + }) + })?; + let total_content_bytes = self + .total_content_bytes + .checked_add(page_content_bytes) + .ok_or_else(|| { + aggregation_error("portable bootstrap aggregation content bytes overflowed") + })?; + if total_content_bytes > limits.max_total_content_bytes { + return Err(aggregation_error( + "portable bootstrap aggregation exceeded its content byte limit", + )); + } + + self.total_changes = total_changes; + self.total_content_bytes = total_content_bytes; + self.completeness.merge(preserved_completeness); + self.source_versions.extend( + page.source_versions + .into_iter() + .map(|source| (source.source_object.remote_id.clone(), source)), + ); + self.observed_changes.extend( + page.observed_changes + .into_iter() + .map(|change| (change.source_object.remote_id.clone(), change)), + ); + self.contents.extend( + page.contents + .into_iter() + .map(|content| (content.artifact_key.clone(), content)), + ); + for projection in page.projections { + self.projection_paths + .insert(projection.logical_path.as_str().to_string()); + self.projections + .insert(projection.artifact_key.clone(), projection); + } + Ok(()) + } + + fn finish( + self, + next_checkpoint: locality_connector::PortableCheckpoint, + ) -> UnpersistedSynchronizationBatch { + let publication_eligible = self.completeness.is_complete(); + UnpersistedSynchronizationBatch { + source_connection_id: self.source_connection_id, + observed_changes: self.observed_changes.into_values().collect(), + source_versions: self.source_versions.into_values().collect(), + contents: self.contents.into_values().collect(), + projections: self.projections.into_values().collect(), + next_checkpoint, + completeness: self.completeness, + publication_eligible, + } + } +} + +fn completeness_without_continuation( + completeness: &PortableCompleteness, +) -> (bool, PortableCompleteness) { + let reasons = completeness.incomplete_reasons(); + let continuation = reasons.contains(&PortableIncompleteReason::CheckpointContinuation); + let mut preserved = if reasons.is_empty() && !completeness.is_complete() { + PortableCompleteness::default() + } else { + PortableCompleteness::complete() + }; + for reason in reasons { + if reason != &PortableIncompleteReason::CheckpointContinuation { + preserved.merge(PortableCompleteness::incomplete(reason.clone())); + } + } + (continuation, preserved) +} + +fn validate_continuation_checkpoint( + current: Option<&locality_connector::PortableCheckpoint>, + next: &locality_connector::PortableCheckpoint, + seen: &mut BTreeSet<(u16, String)>, +) -> LocalityResult<()> { + if next.opaque.is_empty() { + return Err(aggregation_error( + "portable bootstrap continuation returned an empty checkpoint", + )); + } + if current == Some(next) { + return Err(aggregation_error( + "portable bootstrap continuation repeated its checkpoint", + )); + } + if !seen.insert(checkpoint_identity(next)) { + return Err(aggregation_error( + "portable bootstrap continuation formed a checkpoint cycle", + )); + } + Ok(()) +} + +fn checkpoint_identity(checkpoint: &locality_connector::PortableCheckpoint) -> (u16, String) { + (checkpoint.format_version, checkpoint.opaque.clone()) +} + +fn aggregation_error(message: &'static str) -> LocalityError { + LocalityError::InvalidState(message.to_string()) +} + /// Run one incremental synchronization checkpoint through fetch and render. pub fn synchronize_and_project_portable( connector: &C, diff --git a/crates/locality-engine/tests/synchronize_project.rs b/crates/locality-engine/tests/synchronize_project.rs index 3e907453..fe2dc487 100644 --- a/crates/locality-engine/tests/synchronize_project.rs +++ b/crates/locality-engine/tests/synchronize_project.rs @@ -1,4 +1,5 @@ use std::collections::{BTreeMap, BTreeSet}; +use std::sync::Mutex; use locality_connector::{ ApplyPlanRequest, ApplyPlanResult, ApplyUndoRequest, ApplyUndoResult, Connector, @@ -15,7 +16,8 @@ use locality_core::portable::{ LogicalPath, ProjectionFileKind, SourceAction, SourceConnectionId, SourceEdge, SourceObject, }; use locality_engine::synchronize_project::{ - bootstrap_and_project, synchronize_and_project_portable, + BootstrapAggregationLimits, bootstrap_and_project, bootstrap_and_project_to_completion, + synchronize_and_project_portable, }; #[derive(Clone)] @@ -166,6 +168,225 @@ impl Connector for FixtureConnector { } } +#[derive(Clone, Copy, Debug, PartialEq, Eq)] +enum PagedFault { + None, + NonContinuationIncomplete, + EmptyCheckpoint, + RepeatedCheckpoint, + CheckpointCycle, + ChangedConnection, + DuplicateSourceVersion, + DuplicateObservedSource, + DuplicateContentArtifact, + DuplicateProjectionArtifact, + DuplicateLogicalPath, +} + +struct PagedFixtureConnector { + fault: PagedFault, + calls: Mutex>>, +} + +impl PagedFixtureConnector { + fn new(fault: PagedFault) -> Self { + Self { + fault, + calls: Mutex::new(Vec::new()), + } + } + + fn calls(&self) -> Vec> { + self.calls.lock().expect("calls lock").clone() + } +} + +impl Connector for PagedFixtureConnector { + fn kind(&self) -> ConnectorKind { + ConnectorKind("paged-fixture") + } + + fn capabilities(&self) -> ConnectorCapabilities { + ConnectorCapabilities::read_only() + } + + fn supported_push_operations(&self) -> BTreeSet { + BTreeSet::new() + } + + fn enumerate(&self, _request: EnumerateRequest) -> LocalityResult> { + Ok(Vec::new()) + } + + fn bootstrap_portable( + &self, + request: PortableBootstrapRequest, + ) -> LocalityResult { + let checkpoint = request + .checkpoint + .as_ref() + .map(|checkpoint| checkpoint.opaque.clone()); + self.calls + .lock() + .expect("calls lock") + .push(checkpoint.clone()); + + let offset: usize = match checkpoint.as_deref() { + None => 0, + Some("cp1") => 1, + Some("cp2") => 2, + Some(_) => { + return Err(locality_core::LocalityError::InvalidState( + "fixture checkpoint is invalid".to_string(), + )); + } + }; + let requested = usize::try_from(request.max_changes).expect("u32 fits usize"); + if requested == 0 { + return Err(locality_core::LocalityError::InvalidState( + "fixture page size is zero".to_string(), + )); + } + let end = offset.saturating_add(requested).min(3); + let mut changes = (offset..end) + .map(|index| paged_change(&request.source_connection_id, index, self.fault)) + .collect::>(); + if self.fault == PagedFault::DuplicateObservedSource && offset == 0 { + changes[0].requires_fetch = false; + changes[0].source_object.deleted = true; + } + + let has_more = end < 3; + let mut next_opaque = match end { + 0 => "cp0", + 1 => "cp1", + 2 => "cp2", + _ => "done", + } + .to_string(); + let mut completeness = if has_more { + PortableCompleteness::incomplete(PortableIncompleteReason::CheckpointContinuation) + } else { + PortableCompleteness::complete() + }; + if self.fault == PagedFault::NonContinuationIncomplete && offset == 0 { + completeness.merge(PortableCompleteness::incomplete( + PortableIncompleteReason::ConnectorLimitation { + code: "fixture_gap".to_string(), + remote_id: Some(RemoteId::new("provider-value")), + }, + )); + } + match self.fault { + PagedFault::EmptyCheckpoint if offset == 0 => next_opaque.clear(), + PagedFault::RepeatedCheckpoint if offset == 1 => { + next_opaque = "cp1".to_string(); + completeness = PortableCompleteness::incomplete( + PortableIncompleteReason::CheckpointContinuation, + ); + } + PagedFault::CheckpointCycle if offset == 2 => { + next_opaque = "cp1".to_string(); + completeness = PortableCompleteness::incomplete( + PortableIncompleteReason::CheckpointContinuation, + ); + } + _ => {} + } + + Ok(PortableChangeBatch { + changes, + next_checkpoint: PortableCheckpoint { + format_version: 1, + opaque: next_opaque, + }, + completeness, + }) + } + + fn sync_portable(&self, request: PortableSyncRequest) -> LocalityResult { + self.bootstrap_portable(PortableBootstrapRequest { + source_connection_id: request.source_connection_id, + scope: request.scope, + checkpoint: Some(request.checkpoint), + max_changes: request.max_changes, + }) + } + + fn fetch_portable(&self, request: PortableFetchRequest) -> LocalityResult { + Ok(PortableFetchResult { + native: NativeEntity { + remote_id: request.remote_id.clone(), + kind: "fixture_page".to_string(), + raw: format!("native:{}", request.remote_id.as_str()).into_bytes(), + }, + provider_version: Some(format!("v1:{}", request.remote_id.as_str())), + completeness: PortableCompleteness::complete(), + }) + } + + fn render_portable( + &self, + request: &PortableRenderRequest, + ) -> LocalityResult { + let remote_id = request.native.remote_id.as_str(); + let canonical_key = if self.fault == PagedFault::DuplicateContentArtifact { + "paged:shared:canonical:v1".to_string() + } else { + format!("paged:{remote_id}:canonical:v1") + }; + let projection_key = if self.fault == PagedFault::DuplicateProjectionArtifact { + "paged:shared:projection:v1".to_string() + } else { + format!("paged:{remote_id}:projection:v1") + }; + let canonical = PortableContentArtifact { + artifact_key: PortableArtifactKey::new(canonical_key), + media_type: "text/markdown".to_string(), + body: format!("canonical:{remote_id}\n").into_bytes(), + }; + Ok(PortableRenderResult { + canonical: canonical.clone(), + projections: vec![PortableProjectionArtifact { + artifact: PortableContentArtifact { + artifact_key: PortableArtifactKey::new(projection_key), + media_type: "text/markdown".to_string(), + body: canonical.body, + }, + logical_path: request.logical_path.clone(), + file_kind: ProjectionFileKind::Markdown, + format_version: request.format_version, + supported_actions: BTreeSet::from([SourceAction::Read]), + }], + completeness: PortableCompleteness::complete(), + }) + } + + fn fetch(&self, _request: FetchRequest) -> LocalityResult { + unreachable!("portable engine uses fetch_portable") + } + + fn render(&self, _entity: &NativeEntity) -> LocalityResult { + unreachable!("portable engine uses render_portable") + } + + fn parse(&self, _document: &CanonicalDocument) -> LocalityResult { + unreachable!("not used") + } + + fn check_concurrency(&self, _request: ApplyPlanRequest<'_>) -> LocalityResult<()> { + unreachable!("not used") + } + + fn apply(&self, _request: ApplyPlanRequest<'_>) -> LocalityResult { + unreachable!("not used") + } + + fn apply_undo(&self, _request: ApplyUndoRequest<'_>) -> LocalityResult { + unreachable!("not used") + } +} + #[test] fn deterministic_retry_returns_identical_unpersisted_candidates_and_hashes() { let connector = FixtureConnector::complete(); @@ -283,6 +504,234 @@ fn mixed_owning_root_provenance_fails_closed() { assert!(error.to_string().contains("ambiguous owning-root")); } +#[test] +fn paginated_bootstrap_matches_the_exact_one_shot_candidate() { + let paginated_connector = PagedFixtureConnector::new(PagedFault::None); + let paginated = bootstrap_and_project_to_completion( + &paginated_connector, + paged_request(1), + 1, + generous_aggregation_limits(), + ) + .expect("paginated aggregate"); + let one_shot_connector = PagedFixtureConnector::new(PagedFault::None); + let one_shot = + bootstrap_and_project(&one_shot_connector, paged_request(100), 1).expect("one shot"); + + assert_eq!(paginated, one_shot); + assert_eq!( + paginated_connector.calls(), + vec![None, Some("cp1".to_string()), Some("cp2".to_string())] + ); + assert!(paginated.is_publication_eligible()); + paginated.require_complete().expect("complete aggregate"); +} + +#[test] +fn paginated_bootstrap_begins_with_the_callers_checkpoint() { + let checkpoint = PortableCheckpoint { + format_version: 1, + opaque: "cp1".to_string(), + }; + let mut aggregate_request = paged_request(1); + aggregate_request.checkpoint = Some(checkpoint.clone()); + let connector = PagedFixtureConnector::new(PagedFault::None); + let aggregate = bootstrap_and_project_to_completion( + &connector, + aggregate_request, + 1, + generous_aggregation_limits(), + ) + .expect("aggregate from caller checkpoint"); + let mut one_shot_request = paged_request(100); + one_shot_request.checkpoint = Some(checkpoint); + let one_shot = bootstrap_and_project( + &PagedFixtureConnector::new(PagedFault::None), + one_shot_request, + 1, + ) + .expect("one shot from caller checkpoint"); + + assert_eq!(aggregate, one_shot); + assert_eq!( + connector.calls(), + vec![Some("cp1".to_string()), Some("cp2".to_string())] + ); +} + +#[test] +fn pagination_removes_only_continuation_incompleteness() { + let connector = PagedFixtureConnector::new(PagedFault::NonContinuationIncomplete); + let aggregate = bootstrap_and_project_to_completion( + &connector, + paged_request(1), + 1, + generous_aggregation_limits(), + ) + .expect("incomplete aggregate"); + + assert!(!aggregate.is_publication_eligible()); + assert_eq!( + aggregate.completeness.incomplete_reasons(), + [PortableIncompleteReason::ConnectorLimitation { + code: "fixture_gap".to_string(), + remote_id: Some(RemoteId::new("provider-value")), + }] + ); +} + +#[test] +fn continuation_checkpoints_must_be_nonempty_changed_and_acyclic() { + for (fault, expected) in [ + (PagedFault::EmptyCheckpoint, "empty checkpoint"), + (PagedFault::RepeatedCheckpoint, "repeated its checkpoint"), + (PagedFault::CheckpointCycle, "checkpoint cycle"), + ] { + let connector = PagedFixtureConnector::new(fault); + let error = bootstrap_and_project_to_completion( + &connector, + paged_request(1), + 1, + generous_aggregation_limits(), + ) + .expect_err("unsafe checkpoint must fail"); + assert!(error.to_string().contains(expected), "{fault:?}: {error}"); + } +} + +#[test] +fn aggregation_maps_changed_connection_failures_to_a_static_error() { + let connector = PagedFixtureConnector::new(PagedFault::ChangedConnection); + let error = bootstrap_and_project_to_completion( + &connector, + paged_request(1), + 1, + generous_aggregation_limits(), + ) + .expect_err("changed connection must fail"); + + assert_eq!( + error.to_string(), + "invalid state: portable bootstrap aggregation page failed" + ); + assert!(!error.to_string().contains("provider-changed-connection")); +} + +#[test] +fn aggregation_rejects_every_cross_page_identity_collision() { + for (fault, expected) in [ + ( + PagedFault::DuplicateSourceVersion, + "repeated a source version", + ), + ( + PagedFault::DuplicateObservedSource, + "repeated an observed source", + ), + ( + PagedFault::DuplicateContentArtifact, + "repeated a content artifact", + ), + ( + PagedFault::DuplicateProjectionArtifact, + "repeated a projection artifact", + ), + (PagedFault::DuplicateLogicalPath, "repeated a logical path"), + ] { + let connector = PagedFixtureConnector::new(fault); + let error = bootstrap_and_project_to_completion( + &connector, + paged_request(1), + 1, + generous_aggregation_limits(), + ) + .expect_err("cross-page collision must fail"); + assert!(error.to_string().contains(expected), "{fault:?}: {error}"); + } +} + +#[test] +fn aggregation_limits_are_nonzero_and_enforced_before_growth() { + for limits in [ + BootstrapAggregationLimits { + max_checkpoints: 0, + ..generous_aggregation_limits() + }, + BootstrapAggregationLimits { + max_total_changes: 0, + ..generous_aggregation_limits() + }, + BootstrapAggregationLimits { + max_total_content_bytes: 0, + ..generous_aggregation_limits() + }, + ] { + let connector = PagedFixtureConnector::new(PagedFault::None); + let error = bootstrap_and_project_to_completion(&connector, paged_request(1), 1, limits) + .expect_err("zero limit must fail"); + assert!(error.to_string().contains("limits must be nonzero")); + assert!(connector.calls().is_empty()); + } + + let checkpoint_connector = PagedFixtureConnector::new(PagedFault::None); + let checkpoint_error = bootstrap_and_project_to_completion( + &checkpoint_connector, + paged_request(1), + 1, + BootstrapAggregationLimits { + max_checkpoints: 2, + ..generous_aggregation_limits() + }, + ) + .expect_err("checkpoint bound"); + assert!(checkpoint_error.to_string().contains("checkpoint limit")); + assert_eq!(checkpoint_connector.calls().len(), 2); + + let change_connector = PagedFixtureConnector::new(PagedFault::None); + let change_error = bootstrap_and_project_to_completion( + &change_connector, + paged_request(1), + 1, + BootstrapAggregationLimits { + max_total_changes: 2, + ..generous_aggregation_limits() + }, + ) + .expect_err("change bound"); + assert!(change_error.to_string().contains("change limit")); + + let content_connector = PagedFixtureConnector::new(PagedFault::None); + let content_error = bootstrap_and_project_to_completion( + &content_connector, + paged_request(1), + 1, + BootstrapAggregationLimits { + max_total_content_bytes: 1, + ..generous_aggregation_limits() + }, + ) + .expect_err("content bound"); + assert!(content_error.to_string().contains("content byte limit")); +} + +#[test] +fn aggregation_preserves_direct_single_page_behavior() { + let direct_connector = PagedFixtureConnector::new(PagedFault::None); + let direct = + bootstrap_and_project(&direct_connector, paged_request(100), 1).expect("direct batch"); + let aggregate_connector = PagedFixtureConnector::new(PagedFault::None); + let aggregate = bootstrap_and_project_to_completion( + &aggregate_connector, + paged_request(100), + 1, + generous_aggregation_limits(), + ) + .expect("single-page aggregate"); + + assert_eq!(aggregate, direct); + assert_eq!(aggregate_connector.calls(), vec![None]); +} + fn request() -> PortableBootstrapRequest { PortableBootstrapRequest { source_connection_id: SourceConnectionId::new("source-fixture"), @@ -292,6 +741,47 @@ fn request() -> PortableBootstrapRequest { } } +fn paged_request(max_changes: u32) -> PortableBootstrapRequest { + PortableBootstrapRequest { + source_connection_id: SourceConnectionId::new("paged-source"), + scope: locality_connector::PortableSourceScope::explicit_roots([RemoteId::new("root")]), + checkpoint: None, + max_changes, + } +} + +fn generous_aggregation_limits() -> BootstrapAggregationLimits { + BootstrapAggregationLimits { + max_checkpoints: 10, + max_total_changes: 100, + max_total_content_bytes: 1_000_000, + } +} + +fn paged_change( + requested_connection: &SourceConnectionId, + index: usize, + fault: PagedFault, +) -> PortableSourceChange { + let connection = if fault == PagedFault::ChangedConnection { + SourceConnectionId::new("provider-changed-connection") + } else { + requested_connection.clone() + }; + let remote_id = match (fault, index) { + (PagedFault::DuplicateSourceVersion | PagedFault::DuplicateObservedSource, 1) => "page-a", + (_, 0) => "page-a", + (_, 1) => "page-b", + _ => "page-c", + }; + let path = if fault == PagedFault::DuplicateLogicalPath { + "Shared/page.md".to_string() + } else { + format!("{remote_id}/page.md") + }; + change(&connection, remote_id, &path) +} + fn change( source_connection_id: &SourceConnectionId, remote_id: &str, From f70c91fc49d3923d41166eaca3aeb4fec2603820 Mon Sep 17 00:00:00 2001 From: misrasaurabh1 Date: Wed, 22 Jul 2026 15:55:01 -0700 Subject: [PATCH 3/3] Preserve safe portable Notion external media references --- crates/locality-notion/src/dto.rs | 5 +- crates/locality-notion/src/media.rs | 32 ++ crates/locality-notion/src/portable.rs | 144 ++++++-- crates/locality-notion/tests/fetch_render.rs | 362 ++++++++++++++++++- 4 files changed, 513 insertions(+), 30 deletions(-) diff --git a/crates/locality-notion/src/dto.rs b/crates/locality-notion/src/dto.rs index 1f9b11f3..89d64974 100644 --- a/crates/locality-notion/src/dto.rs +++ b/crates/locality-notion/src/dto.rs @@ -16,8 +16,9 @@ pub struct NotionPageBundle { /// Versioned native page payload used only by opt-in portable hosted-media capture. /// /// `page` contains a sanitized Notion bundle: hosted URLs have no query or -/// fragment, expiries are removed, and uncaptured media URLs are empty. Captured -/// bytes are keyed by stable block identity rather than by signed provider URL. +/// fragment, expiries are removed, uncaptured hosted URLs are empty, and safe +/// external HTTPS references retain their exact spelling. Captured bytes are +/// keyed by stable block identity rather than by signed provider URL. #[derive(Clone, Debug, PartialEq, Eq, Serialize, Deserialize)] pub struct NotionPortablePageBundleV1 { pub format_version: u16, diff --git a/crates/locality-notion/src/media.rs b/crates/locality-notion/src/media.rs index 59916bc7..f0b3ef55 100644 --- a/crates/locality-notion/src/media.rs +++ b/crates/locality-notion/src/media.rs @@ -33,6 +33,7 @@ pub const PORTABLE_MEDIA_MAX_ASSETS: usize = 128; pub const PORTABLE_MEDIA_MAX_ASSET_BYTES: usize = 20 * 1024 * 1024; pub const PORTABLE_MEDIA_MAX_AGGREGATE_BYTES: usize = 100 * 1024 * 1024; const PORTABLE_MEDIA_READ_BUFFER_BYTES: usize = 64 * 1024; +const PORTABLE_EXTERNAL_MEDIA_MAX_URL_BYTES: usize = 8 * 1024; const PORTABLE_MEDIA_MAX_REDIRECTS: usize = 3; const PORTABLE_MEDIA_CONNECT_TIMEOUT: Duration = Duration::from_secs(5); const PORTABLE_MEDIA_REQUEST_TIMEOUT: Duration = Duration::from_secs(30); @@ -313,6 +314,37 @@ pub(crate) fn sanitize_portable_hosted_media_url(url: &str) -> LocalityResult LocalityResult<()> { + if url.is_empty() + || url.len() > PORTABLE_EXTERNAL_MEDIA_MAX_URL_BYTES + || url + .chars() + .any(|character| character.is_ascii_control() || character.is_whitespace()) + { + return Err(LocalityError::InvalidState( + "portable external media URL is not raw-safe".to_string(), + )); + } + let parsed = reqwest::Url::parse(url).map_err(|_| { + LocalityError::InvalidState("portable external media URL is invalid".to_string()) + })?; + if parsed.scheme() != "https" + || parsed.host_str().is_none() + || !parsed.username().is_empty() + || parsed.password().is_some() + { + return Err(LocalityError::InvalidState( + "portable external media URL violates the HTTPS reference policy".to_string(), + )); + } + Ok(()) +} + pub(crate) fn portable_media_expired(expiry_time: &str) -> LocalityResult { let expiry = parse_rfc3339_utc_seconds(expiry_time).ok_or_else(|| { LocalityError::InvalidState("portable media expiry is invalid".to_string()) diff --git a/crates/locality-notion/src/portable.rs b/crates/locality-notion/src/portable.rs index 703916dd..eacfde6b 100644 --- a/crates/locality-notion/src/portable.rs +++ b/crates/locality-notion/src/portable.rs @@ -36,6 +36,7 @@ use crate::media::{ PORTABLE_MEDIA_MAX_AGGREGATE_BYTES, PORTABLE_MEDIA_MAX_ASSET_BYTES, PORTABLE_MEDIA_MAX_ASSETS, PortableMediaCaptureFetcher, PortableMediaCapturePolicy, default_portable_media_fetcher, portable_media_expired, sanitize_portable_hosted_media_url, sanitize_portable_media_type, + validate_portable_external_media_url, }; use crate::projection::enumerate_explicit_root_trees; use crate::render::{RenderOptions, render_native_entity, render_native_entity_with_options}; @@ -348,8 +349,34 @@ impl<'a> PortableMediaCaptureState<'a> { payload: &mut FileBlockDto, ) -> LocalityResult<()> { let external_present = payload.external.is_some(); + let hosted_present = payload.file.is_some(); + if external_present && hosted_present { + if let Some(external) = payload.external.as_mut() { + external.url.clear(); + } + if let Some(hosted) = payload.file.as_mut() { + hosted.url.clear(); + hosted.expiry_time = None; + } + if !self.limit_exceeded { + self.record_incomplete(block_id, kind, "ambiguous_file_source"); + } + return Ok(()); + } if let Some(external) = payload.external.as_mut() { + if payload.kind != "external" { + return Err(LocalityError::InvalidState( + "Notion portable media payload type does not match its source".to_string(), + )); + } + if validate_portable_external_media_url(&external.url).is_ok() { + return Ok(()); + } external.url.clear(); + if !self.limit_exceeded { + self.record_incomplete(block_id, kind, "invalid_external_media"); + } + return Ok(()); } if self.limit_exceeded { if let Some(hosted) = payload.file.as_mut() { @@ -359,24 +386,22 @@ impl<'a> PortableMediaCaptureState<'a> { return Ok(()); } let Some(hosted) = payload.file.as_mut() else { - self.record_incomplete( - block_id, - kind, - if external_present { - "external_media" - } else { - "missing_file" - }, - ); + if !matches!(payload.kind.as_str(), "external" | "file") { + return Err(LocalityError::InvalidState( + "Notion portable media payload type does not match its source".to_string(), + )); + } + self.record_incomplete(block_id, kind, "missing_file"); return Ok(()); }; + if payload.kind != "file" { + return Err(LocalityError::InvalidState( + "Notion portable media payload type does not match its source".to_string(), + )); + } let original_url = hosted.url.clone(); let expiry_time = hosted.expiry_time.take(); hosted.url.clear(); - if external_present { - self.record_incomplete(block_id, kind, "ambiguous_file_source"); - return Ok(()); - } if original_url.is_empty() { self.record_incomplete(block_id, kind, "unavailable_hosted_media"); return Ok(()); @@ -513,6 +538,21 @@ fn is_media_kind(kind: &str) -> bool { matches!(kind, "image" | "video" | "file" | "pdf" | "audio") } +fn validate_portable_media_payload_source(payload: &FileBlockDto) -> LocalityResult<()> { + let source_matches = match (payload.external.is_some(), payload.file.is_some()) { + (true, true) => true, + (true, false) => payload.kind == "external", + (false, true) => payload.kind == "file", + (false, false) => matches!(payload.kind.as_str(), "external" | "file"), + }; + if !source_matches { + return Err(LocalityError::InvalidState( + "Notion portable media payload type does not match its source".to_string(), + )); + } + Ok(()) +} + fn portable_media_asset_count(bundle: &NotionPageBundle) -> LocalityResult { fn count_blocks( trees: &[BlockTreeDto], @@ -527,11 +567,24 @@ fn portable_media_asset_count(bundle: &NotionPageBundle) -> LocalityResult>(); + let mut external_block_ids = portable_external_media_block_ids(&bundle.page.blocks); let page_native = NativeEntity { remote_id: request.native.remote_id.clone(), kind: "notion_page".to_string(), @@ -976,6 +1030,9 @@ fn render_portable_media_page( let mut projected_paths = BTreeSet::new(); for rendered_asset in rendered.media_assets { let Some(captured) = captured_by_block.remove(&rendered_asset.block_id) else { + if external_block_ids.remove(&rendered_asset.block_id) { + continue; + } return Err(LocalityError::InvalidState( "Notion portable media render produced an uncaptured asset".to_string(), )); @@ -1026,6 +1083,11 @@ fn render_portable_media_page( "Notion portable media native payload contains an unrendered asset".to_string(), )); } + if !external_block_ids.is_empty() { + return Err(LocalityError::InvalidState( + "Notion portable external media did not render as a reference".to_string(), + )); + } Ok(PortableRenderResult { canonical, @@ -1082,19 +1144,21 @@ fn validate_portable_media_bundle(bundle: &NotionPortablePageBundleV1) -> Locali )); } for (block_id, (kind, payload)) in &media_blocks { + validate_portable_media_payload_source(payload)?; if payload - .external + .file .as_ref() - .is_some_and(|file| !file.url.is_empty()) - || payload - .file - .as_ref() - .is_some_and(|file| file.expiry_time.is_some()) + .is_some_and(|file| file.expiry_time.is_some()) { return Err(LocalityError::InvalidState( "Notion portable media native payload is not sanitized".to_string(), )); } + let external_url = payload + .external + .as_ref() + .map(|file| file.url.as_str()) + .unwrap_or(""); let hosted_url = payload .file .as_ref() @@ -1122,6 +1186,19 @@ fn validate_portable_media_bundle(bundle: &NotionPortablePageBundleV1) -> Locali )); } } else { + if !external_url.is_empty() { + if payload.file.is_some() || payload.kind != "external" { + return Err(LocalityError::InvalidState( + "Notion portable external media source is ambiguous".to_string(), + )); + } + validate_portable_external_media_url(external_url).map_err(|_| { + LocalityError::InvalidState( + "Notion portable external media URL is not allowed".to_string(), + ) + })?; + continue; + } if !hosted_url.is_empty() { return Err(LocalityError::InvalidState( "Notion portable incomplete media retained a remote URL".to_string(), @@ -1130,7 +1207,7 @@ fn validate_portable_media_bundle(bundle: &NotionPortablePageBundleV1) -> Locali if !limit_exceeded { let code = match (payload.external.is_some(), payload.file.is_some()) { (true, true) => "ambiguous_file_source", - (true, false) => "external_media", + (true, false) => "invalid_external_media", (false, true) => "unavailable_hosted_media", (false, false) => "missing_file", }; @@ -1229,6 +1306,29 @@ fn validate_portable_media_bundle(bundle: &NotionPortablePageBundleV1) -> Locali Ok(()) } +fn portable_external_media_block_ids(trees: &[BlockTreeDto]) -> BTreeSet { + fn collect(trees: &[BlockTreeDto], block_ids: &mut BTreeSet) { + for tree in trees { + if is_media_kind(&tree.block.kind) + && media_payload(&tree.block).is_some_and(|payload| { + payload.file.is_none() + && payload.kind == "external" + && payload.external.as_ref().is_some_and(|external| { + validate_portable_external_media_url(&external.url).is_ok() + }) + }) + { + block_ids.insert(tree.block.id.clone()); + } + collect(&tree.children, block_ids); + } + } + + let mut block_ids = BTreeSet::new(); + collect(trees, &mut block_ids); + block_ids +} + fn insert_expected_incomplete( expected: &mut BTreeMap, block_id: &str, diff --git a/crates/locality-notion/tests/fetch_render.rs b/crates/locality-notion/tests/fetch_render.rs index 18baedf3..38a73bc6 100644 --- a/crates/locality-notion/tests/fetch_render.rs +++ b/crates/locality-notion/tests/fetch_render.rs @@ -3156,6 +3156,191 @@ fn portable_media_default_preserves_native_and_remains_incomplete() { assert_eq!(rendered.projections.len(), 1); } +#[test] +fn portable_external_media_is_an_exact_byte_free_reference_and_never_fetched() { + let page_id = "aaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaa"; + let block_id = "cccccccccccccccccccccccccccccccc"; + let external_url = "https://CDN.Example.com:8443/assets/My%20Image.PNG?download=One#Section"; + let calls = Arc::new(Mutex::new(Vec::new())); + let fetcher = Arc::new(FixturePortableMediaFetcher { + outcomes: BTreeMap::new(), + calls: Arc::clone(&calls), + }); + let block = file_block(block_id, "image", external_url, "Image caption"); + let direct_bundle = locality_notion::dto::NotionPageBundle { + page: page(page_id, "Coverage"), + blocks: vec![BlockTreeDto { + block: block.clone(), + children: Vec::new(), + }], + }; + let direct = locality_notion::render::render_page_bundle_with_options( + &direct_bundle, + &locality_notion::render::RenderOptions::with_page_path("Docs/Coverage/page.md") + .with_local_media_block_ids(Vec::::new()), + ) + .expect("direct external reference render"); + let connector = portable_media_connector(page_id, vec![block]) + .with_portable_media_capture_fetcher(PortableMediaCapturePolicy::HostedPilot, fetcher); + + let fetched = connector + .fetch_portable(portable_fetch_request(page_id)) + .expect("portable external reference fetch"); + assert!(calls.lock().expect("calls").is_empty()); + assert!(fetched.completeness.is_complete()); + let native: NotionPortablePageBundleV1 = + serde_json::from_slice(&fetched.native.raw).expect("portable media native"); + assert!(native.captured_media.is_empty()); + assert!(native.incomplete_media.is_empty()); + assert_eq!( + native.page.blocks[0] + .block + .image + .as_ref() + .and_then(|payload| payload.external.as_ref()) + .map(|external| external.url.as_str()), + Some(external_url) + ); + + let rendered = connector + .render_portable(&PortableRenderRequest { + source_connection_id: SourceConnectionId::new("source-notion"), + logical_path: LogicalPath::new("Docs/Coverage/page.md").expect("path"), + native: fetched.native, + format_version: 1, + }) + .expect("portable external reference render"); + let exact_markdown = concat!( + "---\n", + "loc:\n", + " id: aaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaa\n", + " type: page\n", + " synced_at: \"2026-06-10T00:00:00.000Z\"\n", + " remote_edited_at: \"2026-06-10T00:00:00.000Z\"\n", + "title: \"Coverage\"\n", + "---\n", + "![Image caption](https://CDN.Example.com:8443/assets/My%20Image.PNG?download=One#Section)\n", + ); + assert_eq!(rendered.canonical.body, exact_markdown.as_bytes()); + assert_eq!( + rendered.canonical.body, + render_canonical_markdown(&direct.document).as_bytes() + ); + assert_eq!(rendered.projections.len(), 1); + assert_eq!( + rendered.projections[0].artifact.body, + exact_markdown.as_bytes() + ); + assert!(rendered.completeness.is_complete()); + + let mut invalid_url = native.clone(); + invalid_url.page.blocks[0] + .block + .image + .as_mut() + .expect("image payload") + .external + .as_mut() + .expect("external source") + .url = "http://example.com/not-https.png".to_string(); + let mut mismatched_source = native; + mismatched_source.page.blocks[0] + .block + .image + .as_mut() + .expect("image payload") + .kind = "file".to_string(); + for (case, bundle, expected) in [ + ( + "invalid URL", + invalid_url, + "invalid state: Notion portable external media URL is not allowed", + ), + ( + "mismatched source", + mismatched_source, + "invalid state: Notion portable media payload type does not match its source", + ), + ] { + let error = connector + .render_portable(&portable_render_request( + page_id, + NativeEntity { + remote_id: RemoteId::new(page_id), + kind: "notion_page_portable_media_v1".to_string(), + raw: serde_json::to_vec(&bundle).expect("tampered native"), + }, + )) + .unwrap_err(); + assert_eq!(error.to_string(), expected, "{case}"); + } +} + +#[test] +fn portable_external_media_supports_every_file_like_block_without_capture() { + let page_id = "all-external-page"; + let blocks = ["image", "video", "file", "pdf", "audio"] + .into_iter() + .enumerate() + .map(|(index, kind)| { + file_block( + &format!("external-{kind}"), + kind, + &format!("https://example.com/{kind}?source=external&index={index}"), + &format!("External {kind}"), + ) + }) + .collect::>(); + let calls = Arc::new(Mutex::new(Vec::new())); + let connector = portable_media_connector(page_id, blocks).with_portable_media_capture_fetcher( + PortableMediaCapturePolicy::HostedPilot, + Arc::new(FixturePortableMediaFetcher { + outcomes: BTreeMap::new(), + calls: Arc::clone(&calls), + }), + ); + + let fetched = connector + .fetch_portable(portable_fetch_request(page_id)) + .expect("all external media fetch"); + assert!(calls.lock().expect("calls").is_empty()); + assert!(fetched.completeness.is_complete()); + let native: NotionPortablePageBundleV1 = + serde_json::from_slice(&fetched.native.raw).expect("portable native"); + assert!(native.captured_media.is_empty()); + assert!(native.incomplete_media.is_empty()); + + let rendered = connector + .render_portable(&portable_render_request(page_id, fetched.native)) + .expect("all external media render"); + let expected_body = concat!( + "![External image](https://example.com/image?source=external&index=0)\n\n", + "[External video](https://example.com/video?source=external&index=1)\n\n", + "[External file](https://example.com/file?source=external&index=2)\n\n", + "[External pdf](https://example.com/pdf?source=external&index=3)\n\n", + "[External audio](https://example.com/audio?source=external&index=4)\n", + ); + assert_eq!( + String::from_utf8(rendered.canonical.body).expect("UTF-8 Markdown"), + format!( + concat!( + "---\n", + "loc:\n", + " id: all-external-page\n", + " type: page\n", + " synced_at: \"2026-06-10T00:00:00.000Z\"\n", + " remote_edited_at: \"2026-06-10T00:00:00.000Z\"\n", + "title: \"Coverage\"\n", + "---\n", + "{}" + ), + expected_body + ) + ); + assert_eq!(rendered.projections.len(), 1); + assert!(rendered.completeness.is_complete()); +} + #[test] fn portable_capture_keeps_pages_without_media_byte_exact() { let page_id = "aaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaa"; @@ -3196,7 +3381,6 @@ fn portable_capture_keeps_pages_without_media_byte_exact() { #[test] fn portable_media_denials_are_incomplete_and_never_publish_remote_urls() { let denied = [ - ("external", "https://example.com/external.png"), ("http", "http://secure.notion-static.com/image.png"), ("ip", "https://127.0.0.1/image.png"), ( @@ -3213,11 +3397,7 @@ fn portable_media_denials_are_incomplete_and_never_publish_remote_urls() { for (index, (case, url)) in denied.into_iter().enumerate() { let page_id = format!("page-denied-{index}"); let block_id = format!("block-denied-{index}"); - let block = if case == "external" { - file_block(&block_id, "image", url, "Image") - } else { - hosted_file_block(&block_id, "image", url, None) - }; + let block = hosted_file_block(&block_id, "image", url, None); let connector = portable_media_connector(&page_id, vec![block]) .with_portable_media_capture(PortableMediaCapturePolicy::HostedPilot); let fetched = connector @@ -3240,6 +3420,176 @@ fn portable_media_denials_are_incomplete_and_never_publish_remote_urls() { } } +#[test] +fn portable_external_media_fails_closed_without_fetching_bad_or_ambiguous_sources() { + let invalid = [ + ("empty", ""), + ("malformed", "not a URL"), + ("http", "http://example.com/image.png"), + ("userinfo", "https://user:pass@example.com/image.png"), + ]; + for (index, (case, url)) in invalid.into_iter().enumerate() { + let page_id = format!("external-invalid-{index}"); + let block_id = format!("external-block-{index}"); + let calls = Arc::new(Mutex::new(Vec::new())); + let connector = + portable_media_connector(&page_id, vec![file_block(&block_id, "image", url, "Image")]) + .with_portable_media_capture_fetcher( + PortableMediaCapturePolicy::HostedPilot, + Arc::new(FixturePortableMediaFetcher { + outcomes: BTreeMap::new(), + calls: Arc::clone(&calls), + }), + ); + let fetched = connector + .fetch_portable(portable_fetch_request(&page_id)) + .unwrap_or_else(|error| panic!("{case} external fetch: {error}")); + assert!(calls.lock().expect("calls").is_empty(), "{case}"); + assert!(!fetched.completeness.is_complete(), "{case}"); + let native: NotionPortablePageBundleV1 = + serde_json::from_slice(&fetched.native.raw).expect("portable native"); + assert!(native.captured_media.is_empty(), "{case}"); + assert_eq!( + native.incomplete_media, + vec![NotionPortableIncompleteMediaV1 { + block_id: block_id.clone(), + kind: "image".to_string(), + code: "invalid_external_media".to_string(), + }], + "{case}" + ); + assert_eq!( + native.page.blocks[0] + .block + .image + .as_ref() + .and_then(|payload| payload.external.as_ref()) + .map(|external| external.url.as_str()), + Some(""), + "{case}" + ); + let rendered = connector + .render_portable(&portable_render_request( + &page_id, + NativeEntity { + remote_id: RemoteId::new(&page_id), + kind: "notion_page_portable_media_v1".to_string(), + raw: serde_json::to_vec(&native).expect("portable native"), + }, + )) + .unwrap_or_else(|error| panic!("{case} external render: {error}")); + assert_eq!(rendered.projections.len(), 1, "{case}"); + assert!(!rendered.completeness.is_complete(), "{case}"); + if !url.is_empty() { + assert!( + !String::from_utf8_lossy(&rendered.canonical.body).contains(url), + "{case}" + ); + } + } + + for (case, url) in [ + ( + "control", + "https://example.com/image.png\u{7f}hidden".to_string(), + ), + ( + "oversized", + format!("https://example.com/{}", "a".repeat(8 * 1024)), + ), + ] { + let calls = Arc::new(Mutex::new(Vec::new())); + let connector = portable_media_connector( + case, + vec![file_block("invalid-raw", "image", &url, "Image")], + ) + .with_portable_media_capture_fetcher( + PortableMediaCapturePolicy::HostedPilot, + Arc::new(FixturePortableMediaFetcher { + outcomes: BTreeMap::new(), + calls: Arc::clone(&calls), + }), + ); + let fetched = connector + .fetch_portable(portable_fetch_request(case)) + .expect("raw-unsafe external URL becomes explicitly incomplete"); + assert!(calls.lock().expect("calls").is_empty(), "{case}"); + assert!(!fetched.completeness.is_complete(), "{case}"); + assert!(!String::from_utf8_lossy(&fetched.native.raw).contains(&url)); + } + + let external_url = "https://example.com/public.png"; + let hosted_url = "https://secure.notion-static.com/private.png?X-Amz-Signature=must-not-escape"; + let mut ambiguous = file_block("ambiguous", "image", external_url, "Image"); + ambiguous.image.as_mut().expect("image payload").file = Some(HostedFileDto { + url: hosted_url.to_string(), + expiry_time: Some("2099-01-01T00:00:00.000Z".to_string()), + }); + let calls = Arc::new(Mutex::new(Vec::new())); + let connector = portable_media_connector("ambiguous-page", vec![ambiguous]) + .with_portable_media_capture_fetcher( + PortableMediaCapturePolicy::HostedPilot, + Arc::new(FixturePortableMediaFetcher { + outcomes: BTreeMap::new(), + calls: Arc::clone(&calls), + }), + ); + let fetched = connector + .fetch_portable(portable_fetch_request("ambiguous-page")) + .expect("ambiguous source becomes explicit incomplete"); + assert!(calls.lock().expect("calls").is_empty()); + assert!(!fetched.completeness.is_complete()); + let native: NotionPortablePageBundleV1 = + serde_json::from_slice(&fetched.native.raw).expect("ambiguous native"); + assert!(native.captured_media.is_empty()); + assert_eq!( + native.incomplete_media, + vec![NotionPortableIncompleteMediaV1 { + block_id: "ambiguous".to_string(), + kind: "image".to_string(), + code: "ambiguous_file_source".to_string(), + }] + ); + let raw = String::from_utf8_lossy(&fetched.native.raw); + assert!(!raw.contains(external_url)); + assert!(!raw.contains(hosted_url)); + + for (case, mut block) in [ + ( + "external-as-file", + file_block("mismatch-external", "image", external_url, "Image"), + ), + ( + "hosted-as-external", + hosted_file_block("mismatch-hosted", "image", hosted_url, None), + ), + ] { + let payload = block.image.as_mut().expect("image payload"); + payload.kind = if payload.kind == "external" { + "file".to_string() + } else { + "external".to_string() + }; + let calls = Arc::new(Mutex::new(Vec::new())); + let connector = portable_media_connector(case, vec![block]) + .with_portable_media_capture_fetcher( + PortableMediaCapturePolicy::HostedPilot, + Arc::new(FixturePortableMediaFetcher { + outcomes: BTreeMap::new(), + calls: Arc::clone(&calls), + }), + ); + assert_eq!( + connector + .fetch_portable(portable_fetch_request(case)) + .expect_err("payload mismatch must fail") + .to_string(), + "invalid state: Notion portable media payload type does not match its source" + ); + assert!(calls.lock().expect("calls").is_empty(), "{case}"); + } +} + #[test] fn portable_media_expired_failed_and_oversized_captures_are_redacted() { let cases = [