Skip to content
Closed
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
122 changes: 102 additions & 20 deletions home-mixer/candidate_hydrators/subscription_hydrator.rs
Original file line number Diff line number Diff line change
Expand Up @@ -6,9 +6,15 @@ use tonic::async_trait;
use xai_candidate_pipeline::component_library::utils::{default_quick_cache, QuickCache};
use xai_candidate_pipeline::hydrator::{CacheStore, CachedHydrator};

#[derive(Clone, Debug, PartialEq, Eq)]
pub enum SubscriptionAuthorHydration {
Resolved(Option<u64>),
Failed,
}

pub struct SubscriptionHydrator {
pub tes_client: Arc<dyn TESClient + Send + Sync>,
pub cache: QuickCache<u64, Option<u64>>,
pub cache: QuickCache<u64, SubscriptionAuthorHydration>,
}

impl SubscriptionHydrator {
Expand All @@ -21,7 +27,7 @@ impl SubscriptionHydrator {
#[async_trait]
impl CachedHydrator<ScoredPostsQuery, PostCandidate> for SubscriptionHydrator {
type CacheKey = u64;
type CacheValue = Option<u64>;
type CacheValue = SubscriptionAuthorHydration;

fn enable(&self, query: &ScoredPostsQuery) -> bool {
!query.has_cached_posts
Expand All @@ -35,14 +41,11 @@ impl CachedHydrator<ScoredPostsQuery, PostCandidate> for SubscriptionHydrator {
}

fn cache_value(&self, hydrated: &PostCandidate) -> Self::CacheValue {
hydrated.subscription_author_id
cache_value(hydrated)
}

fn hydrate_from_cache(&self, value: Self::CacheValue) -> PostCandidate {
PostCandidate {
subscription_author_id: value,
..Default::default()
}
hydrate_from_cache(value)
}

async fn hydrate_from_client(
Expand All @@ -58,25 +61,104 @@ impl CachedHydrator<ScoredPostsQuery, PostCandidate> for SubscriptionHydrator {

let mut hydrated_candidates = Vec::with_capacity(candidates.len());
for tweet_id in tweet_ids {
let post_features = post_features.get(&tweet_id);
let hydrated = match post_features {
Some(Ok(value)) => Ok(PostCandidate {
subscription_author_id: *value,
..Default::default()
}),
None => Err(format!(
"Missing subscription author id for tweet_id={}",
tweet_id
)),
Some(Err(err)) => Err(err.to_string()),
};
hydrated_candidates.push(hydrated);
hydrated_candidates.push(Ok(hydrate_from_tes(post_features.get(&tweet_id))));
}

hydrated_candidates
}

fn update(&self, candidate: &mut PostCandidate, hydrated: PostCandidate) {
candidate.subscription_author_id = hydrated.subscription_author_id;
candidate.subscription_lookup_failed = hydrated.subscription_lookup_failed;
}
}

fn cache_value(hydrated: &PostCandidate) -> SubscriptionAuthorHydration {
if hydrated.subscription_lookup_failed == Some(true) {
SubscriptionAuthorHydration::Failed
} else {
SubscriptionAuthorHydration::Resolved(hydrated.subscription_author_id)
}
}

fn hydrate_from_tes<E>(result: Option<&Result<Option<u64>, E>>) -> PostCandidate {
match result {
Some(Ok(value)) => PostCandidate {
subscription_author_id: *value,
subscription_lookup_failed: Some(false),
..Default::default()
},
None | Some(Err(_)) => PostCandidate {
subscription_lookup_failed: Some(true),
..Default::default()
},
}
}

fn hydrate_from_cache(value: SubscriptionAuthorHydration) -> PostCandidate {
match value {
SubscriptionAuthorHydration::Failed => PostCandidate {
subscription_lookup_failed: Some(true),
..Default::default()
},
SubscriptionAuthorHydration::Resolved(author_id) => PostCandidate {
subscription_author_id: author_id,
subscription_lookup_failed: Some(false),
..Default::default()
},
}
}

#[cfg(test)]
mod tests {
use super::*;

#[test]
fn cache_round_trip_preserves_lookup_failure() {
let failed = PostCandidate {
subscription_lookup_failed: Some(true),
..Default::default()
};
let cached = cache_value(&failed);
assert_eq!(cached, SubscriptionAuthorHydration::Failed);
let restored = hydrate_from_cache(cached);
assert_eq!(restored.subscription_lookup_failed, Some(true));
assert!(restored.subscription_author_id.is_none());
}

#[test]
fn tes_err_or_missing_slot_fails_closed() {
let err: Result<Option<u64>, &str> = Err("tes unavailable");
let failed = hydrate_from_tes(Some(&err));
assert_eq!(failed.subscription_lookup_failed, Some(true));
assert!(failed.subscription_author_id.is_none());

let missing: Option<&Result<Option<u64>, &str>> = None;
let missing = hydrate_from_tes(missing);
assert_eq!(missing.subscription_lookup_failed, Some(true));

let ok: Result<Option<u64>, &str> = Ok(None);
let not_exclusive = hydrate_from_tes(Some(&ok));
assert_eq!(not_exclusive.subscription_lookup_failed, Some(false));
assert!(not_exclusive.subscription_author_id.is_none());

let exclusive: Result<Option<u64>, &str> = Ok(Some(42));
let exclusive = hydrate_from_tes(Some(&exclusive));
assert_eq!(exclusive.subscription_lookup_failed, Some(false));
assert_eq!(exclusive.subscription_author_id, Some(42));
}

#[test]
fn cache_round_trip_preserves_not_exclusive() {
let resolved = PostCandidate {
subscription_author_id: None,
subscription_lookup_failed: Some(false),
..Default::default()
};
let cached = cache_value(&resolved);
assert_eq!(cached, SubscriptionAuthorHydration::Resolved(None));
let restored = hydrate_from_cache(cached);
assert_eq!(restored.subscription_lookup_failed, Some(false));
assert!(restored.subscription_author_id.is_none());
}
}
35 changes: 29 additions & 6 deletions home-mixer/filters/ineligible_subscription_filter.rs
Original file line number Diff line number Diff line change
Expand Up @@ -19,17 +19,27 @@ impl Filter<ScoredPostsQuery, PostCandidate> for IneligibleSubscriptionFilter {
.collect();

let (kept, removed): (Vec<_>, Vec<_>) =
candidates
.into_iter()
.partition(|candidate| match candidate.subscription_author_id {
Some(author_id) => subscribed_user_ids.contains(&author_id),
None => true,
});
candidates.into_iter().partition(|candidate| {
keep_subscription_candidate(candidate, &subscribed_user_ids)
});

FilterResult { kept, removed }
}
}

fn keep_subscription_candidate(
candidate: &PostCandidate,
subscribed_user_ids: &HashSet<u64>,
) -> bool {
if candidate.subscription_lookup_failed == Some(true) {
return false;
}
match candidate.subscription_author_id {
Some(author_id) => subscribed_user_ids.contains(&author_id),
None => true,
}
}

#[cfg(test)]
mod tests {
use super::*;
Expand Down Expand Up @@ -66,4 +76,17 @@ mod tests {
.iter()
.any(|c| c.subscription_author_id == Some(3)));
}

#[test]
fn tes_lookup_failure_drops_even_when_not_stamped_exclusive() {
let failed = PostCandidate {
subscription_author_id: None,
subscription_lookup_failed: Some(true),
..Default::default()
};
let subscribed = HashSet::from([1u64, 2]);
assert!(!keep_subscription_candidate(&failed, &subscribed));
assert!(keep_subscription_candidate(&candidate(None), &subscribed));
assert!(keep_subscription_candidate(&candidate(Some(1)), &subscribed));
}
}
2 changes: 2 additions & 0 deletions home-mixer/models/candidate.rs
Original file line number Diff line number Diff line change
Expand Up @@ -47,6 +47,8 @@ pub struct PostCandidate {
pub visibility_reason: Option<vf::FilteredReason>,
pub drop_ancillary_posts: Option<bool>,
pub subscription_author_id: Option<u64>,
#[serde(default)]
pub subscription_lookup_failed: Option<bool>,
pub tweet_type_metrics: Option<Vec<u8>>,
pub author_blocks_viewer: Option<bool>,
pub quoted_author_blocks_viewer: Option<bool>,
Expand Down
11 changes: 11 additions & 0 deletions visibility-filtering/hydration/batch.rs
Original file line number Diff line number Diff line change
Expand Up @@ -102,6 +102,13 @@ impl<K: Eq + Hash, V> HydrationBatch<K, V> {
}
}

pub(crate) fn is_failed(&self, key: &K) -> bool {
match self.results.get(key) {
Some(hydrated) => hydrated.is_failed(),
None => true,
}
}

pub(crate) fn hydrated(&self, key: &K) -> Option<&Hydrated<V>> {
self.results.get(key)
}
Expand Down Expand Up @@ -202,6 +209,10 @@ mod tests {
);
assert_eq!(batch.get_or_default(&2), 0);
assert_eq!(batch.get_or_default(&3), 0);
assert!(!batch.is_failed(&1));
assert!(!batch.is_failed(&2));
assert!(batch.is_failed(&3));
assert!(batch.is_failed(&99));
}

#[test]
Expand Down
Loading