refactor!: CQRS repository split — v0.3.0
FederationRepository (34 methods) → 4 focused traits:
ActivityRepository (2) — idempotency tracking
FollowRepository (18) — follower/following graph + migration
ActorRepository (6) — keypairs, remote actor cache, announce tracking
BlocklistRepository (8) — domain + actor blocklists
ApObjectHandler (10 methods) → 2 traits:
ApContentReader (3) — get_local_objects_for_user/page, count_local_posts
ApObjectHandler (9) — all inbox callbacks (on_create, on_mention, etc.)
Builder changes from positional args to named setters:
ActivityPubService::builder(base_url)
.activity_repo(arc)
.follow_repo(arc)
.actor_repo(arc)
.blocklist_repo(arc)
.user_repo(arc)
.content_reader(arc)
.object_handler(arc)
.build()
No behaviour changes.
This commit is contained in:
@@ -86,7 +86,7 @@ impl ActivityPubService {
|
||||
|
||||
loop {
|
||||
let page = data
|
||||
.object_handler
|
||||
.content_reader
|
||||
.get_local_objects_page(owner_user_id, before, BATCH_SIZE)
|
||||
.await?;
|
||||
|
||||
|
||||
@@ -33,7 +33,7 @@ impl ActivityPubService {
|
||||
outbox_url: Some(remote_actor.outbox_url.to_string()),
|
||||
};
|
||||
// Save BEFORE delivering — prevents lost state on process restart.
|
||||
data.federation_repo.add_following(local_user_id, remote, &follow_id_str).await?;
|
||||
data.follow_repo.add_following(local_user_id, remote, &follow_id_str).await?;
|
||||
let follow = FollowActivity {
|
||||
id: Url::parse(&follow_id_str)?,
|
||||
kind: Default::default(),
|
||||
@@ -49,12 +49,12 @@ impl ActivityPubService {
|
||||
if actor_url_str.starts_with(&self.base_url) {
|
||||
return self.unfollow_local(local_user_id, actor_url_str, &data).await;
|
||||
}
|
||||
let remote = data.federation_repo.get_remote_actor(actor_url_str).await?
|
||||
let remote = data.actor_repo.get_remote_actor(actor_url_str).await?
|
||||
.ok_or_else(|| anyhow::anyhow!("remote actor not found: {}", actor_url_str))?;
|
||||
let local_actor = get_local_actor(local_user_id, &data).await.map_err(|e| anyhow::anyhow!("{e}"))?;
|
||||
let remote_ap_id = Url::parse(actor_url_str)?;
|
||||
let inbox = Url::parse(&remote.inbox_url)?;
|
||||
let follow_id = data.federation_repo.get_follow_activity_id(local_user_id, actor_url_str).await?
|
||||
let follow_id = data.follow_repo.get_follow_activity_id(local_user_id, actor_url_str).await?
|
||||
.and_then(|id| Url::parse(&id).ok())
|
||||
.unwrap_or_else(|| activity_url(&self.base_url).unwrap_or_else(|_| remote_ap_id.clone()));
|
||||
let follow = FollowActivity { id: follow_id, kind: Default::default(), actor: ObjectId::from(local_actor.ap_id.clone()), object: ObjectId::from(remote_ap_id) };
|
||||
@@ -66,7 +66,7 @@ impl ActivityPubService {
|
||||
};
|
||||
let (json, sends, inboxes) = self.prepare_broadcast(&data, &local_actor, vec![inbox], undo).await?;
|
||||
self.dispatch_deliveries(&data, &local_actor, inboxes, sends, json).await?;
|
||||
data.federation_repo.remove_following(local_user_id, actor_url_str).await?;
|
||||
data.follow_repo.remove_following(local_user_id, actor_url_str).await?;
|
||||
data.object_handler.on_actor_removed(&Url::parse(actor_url_str)?).await?;
|
||||
Ok(())
|
||||
}
|
||||
@@ -74,13 +74,13 @@ impl ActivityPubService {
|
||||
pub async fn accept_follower(&self, local_user_id: uuid::Uuid, remote_actor_url: &str) -> anyhow::Result<()> {
|
||||
let data = self.federation_config.to_request_data();
|
||||
let local_actor = get_local_actor(local_user_id, &data).await.map_err(|e| anyhow::anyhow!("{e}"))?;
|
||||
let remote_actor = data.federation_repo.get_remote_actor(remote_actor_url).await?
|
||||
let remote_actor = data.actor_repo.get_remote_actor(remote_actor_url).await?
|
||||
.ok_or_else(|| anyhow::anyhow!("remote actor not found"))?;
|
||||
let follow_id_str = data.federation_repo.get_follower_follow_activity_id(local_user_id, remote_actor_url).await?
|
||||
let follow_id_str = data.follow_repo.get_follower_follow_activity_id(local_user_id, remote_actor_url).await?
|
||||
.ok_or_else(|| anyhow::anyhow!("follow activity id not found for {}", remote_actor_url))?;
|
||||
let follow = FollowActivity { id: Url::parse(&follow_id_str)?, kind: Default::default(), actor: ObjectId::from(Url::parse(remote_actor_url)?), object: ObjectId::from(local_actor.ap_id.clone()) };
|
||||
let accept = AcceptActivity { id: activity_url(&self.base_url).map_err(|e| anyhow::anyhow!("{e}"))?, kind: Default::default(), actor: ObjectId::from(local_actor.ap_id.clone()), object: follow };
|
||||
data.federation_repo.update_follower_status(local_user_id, remote_actor_url, FollowerStatus::Accepted).await?;
|
||||
data.follow_repo.update_follower_status(local_user_id, remote_actor_url, FollowerStatus::Accepted).await?;
|
||||
let inbox = Url::parse(&remote_actor.inbox_url)?;
|
||||
let (json, sends, inboxes) = self.prepare_broadcast(&data, &local_actor, vec![inbox], accept).await?;
|
||||
self.dispatch_deliveries(&data, &local_actor, inboxes, sends, json).await?;
|
||||
@@ -92,25 +92,25 @@ impl ActivityPubService {
|
||||
pub async fn reject_follower(&self, local_user_id: uuid::Uuid, remote_actor_url: &str) -> anyhow::Result<()> {
|
||||
let data = self.federation_config.to_request_data();
|
||||
let local_actor = get_local_actor(local_user_id, &data).await.map_err(|e| anyhow::anyhow!("{e}"))?;
|
||||
let remote_actor = data.federation_repo.get_remote_actor(remote_actor_url).await?
|
||||
let remote_actor = data.actor_repo.get_remote_actor(remote_actor_url).await?
|
||||
.ok_or_else(|| anyhow::anyhow!("remote actor not found"))?;
|
||||
let follow = FollowActivity { id: activity_url(&self.base_url).map_err(|e| anyhow::anyhow!("{e}"))?, kind: Default::default(), actor: ObjectId::from(Url::parse(remote_actor_url)?), object: ObjectId::from(local_actor.ap_id.clone()) };
|
||||
let reject = RejectActivity { id: activity_url(&self.base_url).map_err(|e| anyhow::anyhow!("{e}"))?, kind: Default::default(), actor: ObjectId::from(local_actor.ap_id.clone()), object: follow };
|
||||
let inbox = Url::parse(&remote_actor.inbox_url)?;
|
||||
let (json, sends, inboxes) = self.prepare_broadcast(&data, &local_actor, vec![inbox], reject).await?;
|
||||
self.dispatch_deliveries(&data, &local_actor, inboxes, sends, json).await?;
|
||||
data.federation_repo.remove_follower(local_user_id, remote_actor_url).await?;
|
||||
data.follow_repo.remove_follower(local_user_id, remote_actor_url).await?;
|
||||
Ok(())
|
||||
}
|
||||
|
||||
pub async fn get_pending_followers(&self, local_user_id: uuid::Uuid) -> anyhow::Result<Vec<RemoteActor>> {
|
||||
let data = self.federation_config.to_request_data();
|
||||
data.federation_repo.get_pending_followers(local_user_id).await
|
||||
data.follow_repo.get_pending_followers(local_user_id).await
|
||||
}
|
||||
|
||||
pub async fn get_accepted_followers(&self, local_user_id: uuid::Uuid) -> anyhow::Result<Vec<RemoteActor>> {
|
||||
let data = self.federation_config.to_request_data();
|
||||
Ok(data.federation_repo.get_followers(local_user_id).await?
|
||||
Ok(data.follow_repo.get_followers(local_user_id).await?
|
||||
.into_iter()
|
||||
.filter(|f| f.status == FollowerStatus::Accepted)
|
||||
.map(|f| f.actor)
|
||||
@@ -119,7 +119,7 @@ impl ActivityPubService {
|
||||
|
||||
pub async fn count_accepted_followers(&self, local_user_id: uuid::Uuid) -> anyhow::Result<usize> {
|
||||
let data = self.federation_config.to_request_data();
|
||||
Ok(data.federation_repo.get_followers(local_user_id).await?
|
||||
Ok(data.follow_repo.get_followers(local_user_id).await?
|
||||
.into_iter()
|
||||
.filter(|f| f.status == FollowerStatus::Accepted)
|
||||
.count())
|
||||
@@ -127,26 +127,26 @@ impl ActivityPubService {
|
||||
|
||||
pub async fn get_following(&self, local_user_id: uuid::Uuid) -> anyhow::Result<Vec<RemoteActor>> {
|
||||
let data = self.federation_config.to_request_data();
|
||||
data.federation_repo.get_following(local_user_id).await
|
||||
data.follow_repo.get_following(local_user_id).await
|
||||
}
|
||||
|
||||
pub async fn count_following(&self, local_user_id: uuid::Uuid) -> anyhow::Result<usize> {
|
||||
let data = self.federation_config.to_request_data();
|
||||
data.federation_repo.count_following(local_user_id).await
|
||||
data.follow_repo.count_following(local_user_id).await
|
||||
}
|
||||
|
||||
pub async fn remove_follower(&self, local_user_id: uuid::Uuid, actor_url: &str) -> anyhow::Result<()> {
|
||||
let data = self.federation_config.to_request_data();
|
||||
data.federation_repo.remove_follower(local_user_id, actor_url).await
|
||||
data.follow_repo.remove_follower(local_user_id, actor_url).await
|
||||
}
|
||||
|
||||
pub async fn block_actor(&self, local_user_id: uuid::Uuid, actor_url: &str) -> anyhow::Result<()> {
|
||||
let data = self.federation_config.to_request_data();
|
||||
data.federation_repo.add_blocked_actor(local_user_id, actor_url).await?;
|
||||
let _ = data.federation_repo.remove_follower(local_user_id, actor_url).await;
|
||||
let _ = data.federation_repo.remove_following(local_user_id, actor_url).await;
|
||||
data.blocklist_repo.add_blocked_actor(local_user_id, actor_url).await?;
|
||||
let _ = data.follow_repo.remove_follower(local_user_id, actor_url).await;
|
||||
let _ = data.follow_repo.remove_following(local_user_id, actor_url).await;
|
||||
let local_actor = get_local_actor(local_user_id, &data).await.map_err(|e| anyhow::anyhow!("{e}"))?;
|
||||
if let Ok(Some(remote_actor)) = data.federation_repo.get_remote_actor(actor_url).await {
|
||||
if let Ok(Some(remote_actor)) = data.actor_repo.get_remote_actor(actor_url).await {
|
||||
let block = crate::activities::BlockActivity {
|
||||
id: activity_url(&self.base_url).map_err(|e| anyhow::anyhow!("{e}"))?,
|
||||
kind: Default::default(),
|
||||
@@ -162,15 +162,15 @@ impl ActivityPubService {
|
||||
|
||||
pub async fn unblock_actor(&self, local_user_id: uuid::Uuid, actor_url: &str) -> anyhow::Result<()> {
|
||||
let data = self.federation_config.to_request_data();
|
||||
data.federation_repo.remove_blocked_actor(local_user_id, actor_url).await
|
||||
data.blocklist_repo.remove_blocked_actor(local_user_id, actor_url).await
|
||||
}
|
||||
|
||||
pub async fn get_blocked_actors(&self, local_user_id: uuid::Uuid) -> anyhow::Result<Vec<RemoteActor>> {
|
||||
let data = self.federation_config.to_request_data();
|
||||
let actor_urls = data.federation_repo.get_blocked_actors(local_user_id).await?;
|
||||
let actor_urls = data.blocklist_repo.get_blocked_actors(local_user_id).await?;
|
||||
let mut actors = Vec::new();
|
||||
for url in actor_urls {
|
||||
let actor = match data.federation_repo.get_remote_actor(&url).await {
|
||||
let actor = match data.actor_repo.get_remote_actor(&url).await {
|
||||
Ok(Some(a)) => a,
|
||||
_ => RemoteActor { url: url.clone(), handle: url.clone(), inbox_url: url.clone(), shared_inbox_url: None, display_name: None, avatar_url: None, outbox_url: None },
|
||||
};
|
||||
@@ -193,7 +193,7 @@ impl ActivityPubService {
|
||||
let follower_actor_url = crate::urls::actor_url(&self.base_url, local_user_id).to_string();
|
||||
let target_actor_url = crate::urls::actor_url(&self.base_url, target.id);
|
||||
let follow_id = activity_url(&self.base_url).map_err(|e| anyhow::anyhow!("{e}"))?.to_string();
|
||||
data.federation_repo.add_follower(target.id, &follower_actor_url, FollowerStatus::Accepted, &follow_id).await?;
|
||||
data.follow_repo.add_follower(target.id, &follower_actor_url, FollowerStatus::Accepted, &follow_id).await?;
|
||||
let target_as_remote = RemoteActor {
|
||||
url: target_actor_url.to_string(),
|
||||
handle: format!("{}@{}", target.username, data.domain),
|
||||
@@ -203,8 +203,8 @@ impl ActivityPubService {
|
||||
avatar_url: None,
|
||||
outbox_url: None,
|
||||
};
|
||||
data.federation_repo.add_following(local_user_id, target_as_remote, &follow_id).await?;
|
||||
data.federation_repo.update_following_status(local_user_id, target_actor_url.as_ref(), FollowingStatus::Accepted).await?;
|
||||
data.follow_repo.add_following(local_user_id, target_as_remote, &follow_id).await?;
|
||||
data.follow_repo.update_following_status(local_user_id, target_actor_url.as_ref(), FollowingStatus::Accepted).await?;
|
||||
tracing::info!(follower = %local_user_id, followee = %target.id, "local follow");
|
||||
Ok(())
|
||||
}
|
||||
@@ -219,8 +219,8 @@ impl ActivityPubService {
|
||||
let target_user_id = crate::urls::extract_user_id_from_url(&target_url)
|
||||
.ok_or_else(|| anyhow::anyhow!("invalid local actor URL: {}", target_actor_url))?;
|
||||
let local_actor_url = crate::urls::actor_url(&self.base_url, local_user_id).to_string();
|
||||
data.federation_repo.remove_follower(target_user_id, &local_actor_url).await?;
|
||||
data.federation_repo.remove_following(local_user_id, target_actor_url).await?;
|
||||
data.follow_repo.remove_follower(target_user_id, &local_actor_url).await?;
|
||||
data.follow_repo.remove_following(local_user_id, target_actor_url).await?;
|
||||
tracing::info!(follower = %local_user_id, followee = %target_user_id, "local unfollow");
|
||||
Ok(())
|
||||
}
|
||||
|
||||
@@ -10,14 +10,18 @@ use url::Url;
|
||||
use crate::{
|
||||
actor_handler::actor_handler,
|
||||
actors::{DbActor, get_local_actor},
|
||||
content::ApObjectHandler,
|
||||
content::{ApContentReader, ApObjectHandler},
|
||||
data::FederationData,
|
||||
federation::ApFederationConfig,
|
||||
followers_handler::{followers_handler, following_handler},
|
||||
inbox::inbox_handler,
|
||||
nodeinfo::{nodeinfo_handler, nodeinfo_well_known_handler},
|
||||
outbox::outbox_handler,
|
||||
repository::{BlockedDomain, FederationRepository},
|
||||
repository::{
|
||||
ActivityRepository, ActorRepository, BlockedDomain, BlocklistRepository,
|
||||
FollowRepository, FollowerStatus, FollowingStatus, RemoteActor,
|
||||
},
|
||||
urls::activity_url,
|
||||
user::ApUserRepository,
|
||||
webfinger::webfinger_handler,
|
||||
};
|
||||
@@ -45,9 +49,13 @@ pub struct ActivityPubService {
|
||||
}
|
||||
|
||||
pub struct ActivityPubServiceBuilder {
|
||||
repo: Arc<dyn FederationRepository>,
|
||||
user_repo: Arc<dyn ApUserRepository>,
|
||||
object_handler: Arc<dyn crate::content::ApObjectHandler>,
|
||||
activity_repo: Option<Arc<dyn ActivityRepository>>,
|
||||
follow_repo: Option<Arc<dyn FollowRepository>>,
|
||||
actor_repo: Option<Arc<dyn ActorRepository>>,
|
||||
blocklist_repo: Option<Arc<dyn BlocklistRepository>>,
|
||||
user_repo: Option<Arc<dyn ApUserRepository>>,
|
||||
content_reader: Option<Arc<dyn ApContentReader>>,
|
||||
object_handler: Option<Arc<dyn ApObjectHandler>>,
|
||||
base_url: String,
|
||||
allow_registration: bool,
|
||||
software_name: String,
|
||||
@@ -58,19 +66,56 @@ pub struct ActivityPubServiceBuilder {
|
||||
}
|
||||
|
||||
impl ActivityPubServiceBuilder {
|
||||
pub fn activity_repo(mut self, v: Arc<dyn ActivityRepository>) -> Self {
|
||||
self.activity_repo = Some(v); self
|
||||
}
|
||||
pub fn follow_repo(mut self, v: Arc<dyn FollowRepository>) -> Self {
|
||||
self.follow_repo = Some(v); self
|
||||
}
|
||||
pub fn actor_repo(mut self, v: Arc<dyn ActorRepository>) -> Self {
|
||||
self.actor_repo = Some(v); self
|
||||
}
|
||||
pub fn blocklist_repo(mut self, v: Arc<dyn BlocklistRepository>) -> Self {
|
||||
self.blocklist_repo = Some(v); self
|
||||
}
|
||||
pub fn user_repo(mut self, v: Arc<dyn ApUserRepository>) -> Self {
|
||||
self.user_repo = Some(v); self
|
||||
}
|
||||
pub fn content_reader(mut self, v: Arc<dyn ApContentReader>) -> Self {
|
||||
self.content_reader = Some(v); self
|
||||
}
|
||||
pub fn object_handler(mut self, v: Arc<dyn ApObjectHandler>) -> Self {
|
||||
self.object_handler = Some(v); self
|
||||
}
|
||||
pub fn allow_registration(mut self, v: bool) -> Self { self.allow_registration = v; self }
|
||||
pub fn software_name(mut self, v: impl Into<String>) -> Self { self.software_name = v.into(); self }
|
||||
pub fn debug(mut self, v: bool) -> Self { self.debug = v; self }
|
||||
pub fn event_publisher(mut self, v: Arc<dyn crate::data::EventPublisher>) -> Self { self.event_publisher = Some(v); self }
|
||||
/// Override max delivery retries (default: `DELIVERY_MAX_ATTEMPTS`).
|
||||
pub fn event_publisher(mut self, v: Arc<dyn crate::data::EventPublisher>) -> Self {
|
||||
self.event_publisher = Some(v); self
|
||||
}
|
||||
pub fn delivery_max_attempts(mut self, v: u32) -> Self { self.delivery_max_attempts = v; self }
|
||||
/// Override initial retry backoff in seconds (default: `DELIVERY_INITIAL_DELAY_SECS`).
|
||||
pub fn delivery_initial_delay_secs(mut self, v: u64) -> Self { self.delivery_initial_delay_secs = v; self }
|
||||
|
||||
pub async fn build(self) -> anyhow::Result<ActivityPubService> {
|
||||
let activity_repo = self.activity_repo
|
||||
.ok_or_else(|| anyhow::anyhow!("activity_repo required — call .activity_repo(arc)"))?;
|
||||
let follow_repo = self.follow_repo
|
||||
.ok_or_else(|| anyhow::anyhow!("follow_repo required — call .follow_repo(arc)"))?;
|
||||
let actor_repo = self.actor_repo
|
||||
.ok_or_else(|| anyhow::anyhow!("actor_repo required — call .actor_repo(arc)"))?;
|
||||
let blocklist_repo = self.blocklist_repo
|
||||
.ok_or_else(|| anyhow::anyhow!("blocklist_repo required — call .blocklist_repo(arc)"))?;
|
||||
let user_repo = self.user_repo
|
||||
.ok_or_else(|| anyhow::anyhow!("user_repo required — call .user_repo(arc)"))?;
|
||||
let content_reader = self.content_reader
|
||||
.ok_or_else(|| anyhow::anyhow!("content_reader required — call .content_reader(arc)"))?;
|
||||
let object_handler = self.object_handler
|
||||
.ok_or_else(|| anyhow::anyhow!("object_handler required — call .object_handler(arc)"))?;
|
||||
let data = FederationData::new(
|
||||
self.repo, self.user_repo, self.object_handler, self.base_url.clone(),
|
||||
self.allow_registration, self.software_name, self.event_publisher,
|
||||
activity_repo, follow_repo, actor_repo, blocklist_repo,
|
||||
user_repo, content_reader, object_handler,
|
||||
self.base_url.clone(), self.allow_registration, self.software_name,
|
||||
self.event_publisher,
|
||||
);
|
||||
let federation_config = ApFederationConfig::new(data, self.debug).await?;
|
||||
Ok(ActivityPubService {
|
||||
@@ -83,16 +128,20 @@ impl ActivityPubServiceBuilder {
|
||||
}
|
||||
|
||||
impl ActivityPubService {
|
||||
pub fn builder(
|
||||
repo: Arc<dyn FederationRepository>,
|
||||
user_repo: Arc<dyn ApUserRepository>,
|
||||
object_handler: Arc<dyn ApObjectHandler>,
|
||||
base_url: impl Into<String>,
|
||||
) -> ActivityPubServiceBuilder {
|
||||
pub fn builder(base_url: impl Into<String>) -> ActivityPubServiceBuilder {
|
||||
ActivityPubServiceBuilder {
|
||||
repo, user_repo, object_handler, base_url: base_url.into(),
|
||||
allow_registration: false, software_name: String::new(),
|
||||
debug: false, event_publisher: None,
|
||||
activity_repo: None,
|
||||
follow_repo: None,
|
||||
actor_repo: None,
|
||||
blocklist_repo: None,
|
||||
user_repo: None,
|
||||
content_reader: None,
|
||||
object_handler: None,
|
||||
base_url: base_url.into(),
|
||||
allow_registration: false,
|
||||
software_name: String::new(),
|
||||
debug: false,
|
||||
event_publisher: None,
|
||||
delivery_max_attempts: DELIVERY_MAX_ATTEMPTS,
|
||||
delivery_initial_delay_secs: DELIVERY_INITIAL_DELAY_SECS,
|
||||
}
|
||||
@@ -133,11 +182,11 @@ impl ActivityPubService {
|
||||
const PAGE_SIZE: usize = 20;
|
||||
let data = self.federation_config.to_request_data();
|
||||
let collection_id = format!("{}/users/{}/followers", self.base_url, user_id);
|
||||
let total = data.federation_repo.count_followers(user_id).await?;
|
||||
let total = data.follow_repo.count_followers(user_id).await?;
|
||||
let obj = if let Some(p) = page {
|
||||
let p = p.max(1);
|
||||
let offset = (p.saturating_sub(1) as usize) * PAGE_SIZE;
|
||||
let followers = data.federation_repo.get_followers_page(user_id, offset as u32, PAGE_SIZE).await?;
|
||||
let followers = data.follow_repo.get_followers_page(user_id, offset as u32, PAGE_SIZE).await?;
|
||||
let has_next = offset + followers.len() < total;
|
||||
let items: Vec<String> = followers.into_iter().map(|f| f.actor.url).collect();
|
||||
let mut obj = serde_json::json!({"@context":AP_CONTEXT,"type":"OrderedCollectionPage","id":format!("{}?page={}",collection_id,p),"partOf":collection_id,"totalItems":total,"orderedItems":items});
|
||||
@@ -154,11 +203,11 @@ impl ActivityPubService {
|
||||
const PAGE_SIZE: usize = 20;
|
||||
let data = self.federation_config.to_request_data();
|
||||
let collection_id = format!("{}/users/{}/following", self.base_url, user_id);
|
||||
let total = data.federation_repo.count_following(user_id).await?;
|
||||
let total = data.follow_repo.count_following(user_id).await?;
|
||||
let obj = if let Some(p) = page {
|
||||
let p = p.max(1);
|
||||
let offset = (p.saturating_sub(1) as usize) * PAGE_SIZE;
|
||||
let following = data.federation_repo.get_following_page(user_id, offset as u32, PAGE_SIZE).await?;
|
||||
let following = data.follow_repo.get_following_page(user_id, offset as u32, PAGE_SIZE).await?;
|
||||
let has_next = offset + following.len() < total;
|
||||
let items: Vec<String> = following.into_iter().map(|a| a.url).collect();
|
||||
let mut obj = serde_json::json!({"@context":AP_CONTEXT,"type":"OrderedCollectionPage","id":format!("{}?page={}",collection_id,p),"partOf":collection_id,"totalItems":total,"orderedItems":items});
|
||||
@@ -172,12 +221,12 @@ impl ActivityPubService {
|
||||
|
||||
pub async fn mark_follower_accepted(&self, user_id: uuid::Uuid, actor_url: &str) -> anyhow::Result<()> {
|
||||
let data = self.federation_config.to_request_data();
|
||||
data.federation_repo.update_follower_status(user_id, actor_url, crate::repository::FollowerStatus::Accepted).await.map_err(|e| anyhow::anyhow!("{e}"))
|
||||
data.follow_repo.update_follower_status(user_id, actor_url, crate::repository::FollowerStatus::Accepted).await.map_err(|e| anyhow::anyhow!("{e}"))
|
||||
}
|
||||
|
||||
pub async fn mark_follower_rejected(&self, user_id: uuid::Uuid, actor_url: &str) -> anyhow::Result<()> {
|
||||
let data = self.federation_config.to_request_data();
|
||||
data.federation_repo.remove_follower(user_id, actor_url).await.map_err(|e| anyhow::anyhow!("{e}"))
|
||||
data.follow_repo.remove_follower(user_id, actor_url).await.map_err(|e| anyhow::anyhow!("{e}"))
|
||||
}
|
||||
|
||||
pub async fn lookup_actor_by_handle(&self, handle: &str) -> anyhow::Result<crate::user::LookedUpActor> {
|
||||
@@ -200,17 +249,17 @@ impl ActivityPubService {
|
||||
|
||||
pub async fn add_blocked_domain(&self, domain: &str, reason: Option<&str>) -> anyhow::Result<()> {
|
||||
let data = self.federation_config.to_request_data();
|
||||
data.federation_repo.add_blocked_domain(domain, reason).await
|
||||
data.blocklist_repo.add_blocked_domain(domain, reason).await
|
||||
}
|
||||
|
||||
pub async fn remove_blocked_domain(&self, domain: &str) -> anyhow::Result<()> {
|
||||
let data = self.federation_config.to_request_data();
|
||||
data.federation_repo.remove_blocked_domain(domain).await
|
||||
data.blocklist_repo.remove_blocked_domain(domain).await
|
||||
}
|
||||
|
||||
pub async fn get_blocked_domains(&self) -> anyhow::Result<Vec<BlockedDomain>> {
|
||||
let data = self.federation_config.to_request_data();
|
||||
data.federation_repo.get_blocked_domains().await
|
||||
data.blocklist_repo.get_blocked_domains().await
|
||||
}
|
||||
|
||||
// ── Private helpers (accessible to child modules via Rust's privacy rules) ─
|
||||
@@ -221,7 +270,7 @@ impl ActivityPubService {
|
||||
local_user_id: uuid::Uuid,
|
||||
) -> anyhow::Result<Option<(DbActor, Vec<Url>)>> {
|
||||
let local_actor = get_local_actor(local_user_id, data).await.map_err(|e| anyhow::anyhow!("{e}"))?;
|
||||
let inbox_strs = data.federation_repo.get_accepted_follower_inboxes(local_user_id).await?;
|
||||
let inbox_strs = data.follow_repo.get_accepted_follower_inboxes(local_user_id).await?;
|
||||
if inbox_strs.is_empty() { return Ok(None); }
|
||||
let inboxes: Vec<Url> = inbox_strs.into_iter().filter_map(|s| {
|
||||
Url::parse(&s).map_err(|e| tracing::warn!(inbox = %s, error = %e, "skipping unparseable inbox URL")).ok()
|
||||
|
||||
Reference in New Issue
Block a user