use api::heddle::api::v1alpha1::CallFailureCode;
pub use api::heddle::api::v1alpha1::{RepoEvent, SubscribeRepoEventsRequest};
use cli_shared::{RemoteTarget, UserConfig};
pub use crate::hosted_runtime::hosted::HostedError;
use crate::hosted_runtime::{HostedAuthMode, HostedClient, HostedSession, ServerStream};
#[derive(Debug, thiserror::Error)]
pub enum RepoEventError {
#[error("invalid hosted server URL: {0}")]
InvalidServer(String),
#[error("repo-event subscriptions require a hosted network server")]
LocalServer,
#[error("failed to configure the hosted repo-event client: {0}")]
Configuration(String),
#[error("failed to connect the hosted repo-event client: {0}")]
Connection(String),
#[error("repo-event subscription was refused: {source}")]
Refused {
#[source]
source: HostedError,
},
#[error(
"repo-event stream disconnected after event {last_event_id}: {source}; reconnect and resume after event {last_event_id}"
)]
Disconnected {
last_event_id: i64,
#[source]
source: HostedError,
},
#[error(
"repo-event stream ended after event {last_event_id}; reconnect and resume after event {last_event_id}"
)]
Ended {
last_event_id: i64,
},
}
impl RepoEventError {
pub fn resume_after_event_id(&self) -> Option<i64> {
match self {
Self::Disconnected { last_event_id, .. } | Self::Ended { last_event_id } => {
Some(*last_event_id)
}
_ => None,
}
}
}
pub struct RepoEventClient {
client: HostedClient,
}
impl RepoEventClient {
pub async fn connect(server: &str) -> Result<Self, RepoEventError> {
let _ = rustls::crypto::ring::default_provider().install_default();
let target = RemoteTarget::parse_native(server).map_err(RepoEventError::InvalidServer)?;
let RemoteTarget::Network { addr, .. } = target else {
return Err(RepoEventError::LocalServer);
};
let server_key = cli_shared::remote::credential_key_from_remote_url(server);
let user_config = UserConfig::load_default()
.map_err(|error| RepoEventError::Configuration(error.to_string()))?;
let session =
HostedSession::build(&user_config, server_key, HostedAuthMode::CredentialFallback)
.map_err(|error| RepoEventError::Configuration(error.to_string()))?;
let client = session
.connect(addr)
.await
.map_err(|error| RepoEventError::Connection(error.to_string()))?;
Ok(Self { client })
}
pub async fn subscribe(
&self,
request: SubscribeRepoEventsRequest,
) -> Result<RepoEventSubscription, RepoEventError> {
let last_event_id = request.after_event_id.max(0);
let stream = self
.client
.routes()
.subscribe_repo_events(&request)
.await
.map_err(|source| RepoEventError::Disconnected {
last_event_id,
source,
})?;
Ok(RepoEventSubscription {
request,
stream,
last_event_id,
})
}
pub async fn close(self) {
self.client.close().await;
}
#[cfg(test)]
fn from_hosted_client(client: HostedClient) -> Self {
Self { client }
}
}
pub struct RepoEventSubscription {
request: SubscribeRepoEventsRequest,
stream: ServerStream<RepoEvent>,
last_event_id: i64,
}
impl RepoEventSubscription {
pub async fn next(&mut self) -> Result<RepoEvent, RepoEventError> {
match self.stream.next().await {
Ok(Some(event)) => {
self.last_event_id = self.last_event_id.max(event.event_id);
Ok(event)
}
Ok(None) => Err(RepoEventError::Ended {
last_event_id: self.last_event_id,
}),
Err(source) if is_refusal(&source) => Err(RepoEventError::Refused { source }),
Err(source) => Err(RepoEventError::Disconnected {
last_event_id: self.last_event_id,
source,
}),
}
}
pub fn last_event_id(&self) -> i64 {
self.last_event_id
}
pub fn resume_request(&self) -> SubscribeRepoEventsRequest {
let mut request = self.request.clone();
request.after_event_id = self.last_event_id;
request
}
}
fn is_refusal(error: &HostedError) -> bool {
matches!(
error,
HostedError::Call {
code: CallFailureCode::Unauthenticated
| CallFailureCode::PermissionDenied
| CallFailureCode::NotFound,
..
}
)
}
#[cfg(test)]
#[path = "repo_events_tests.rs"]
mod tests;