/
codenik
/
codenik-tools
Обзор
Документация
Войти
/
codenik
/
codenik-tools
Код
Запросы
0
Задачи
Вики
Пакеты
0
Релизы
0
CI/CD
Аналитика
main
crates/plugins/api/sourcecraft/src/client.rs
629 строк
23 KB
Codenik Wizard
feat(sourcecraft): impl Repo/Issue/MergeRequest providers
09 май 2026, 00:32
09 май 2026, 00:32
3d48c2c
Код
Авторство
О чём код?
use async_trait::async_trait; use reqwest::Client; use reqwest::header::{AUTHORIZATION, HeaderMap, HeaderValue, USER_AGENT}; use secrecy::{ExposeSecret, SecretString}; use serde::de::DeserializeOwned; use codenik_core::Commit; use codenik_core::{ Branch, Comment, CommitsOpts, CreateIssueInput, Error, Issue, IssueFilter, IssueKey, IssueProvider, MergeRequest, MergeRequestProvider, MrFilter, Page, Provider, Repo, RepoFilter, RepoProvider, Result, }; use crate::types::{ ApiBranchesResponse, ApiCommentRecord, ApiCommentsResponse, ApiIssue, ApiIssuesResponse, ApiPullRequest, ApiPullRequestsResponse, ApiRepositoriesResponse, ApiRepository, }; const DEFAULT_BASE_URL: &str = "https://api.sourcecraft.tech"; const PROVIDER_NAME: &str = "sourcecraft"; pub struct SourceCraftClient { base_url: String, http: Client, } impl SourceCraftClient { pub fn new(token: SecretString) -> Result<Self> { Self::with_base_url(token, DEFAULT_BASE_URL) } pub fn with_base_url(token: SecretString, base_url: impl Into<String>) -> Result<Self> { let mut headers = HeaderMap::new(); let auth = HeaderValue::from_str(&format!("Bearer {}", token.expose_secret())) .map_err(|e| Error::Config(format!("invalid token header: {e}")))?; headers.insert(AUTHORIZATION, auth); headers.insert(USER_AGENT, HeaderValue::from_static("codenik-tools")); let http = Client::builder() .default_headers(headers) .build() .map_err(|e| Error::Transport(e.to_string()))?; Ok(Self { base_url: base_url.into().trim_end_matches('/').to_string(), http, }) } pub fn base_url(&self) -> &str { &self.base_url } pub(crate) async fn get_json<T: DeserializeOwned>( &self, path: &str, query: &[(String, String)], ) -> Result<T> { let url = format!("{}{}", self.base_url, path); let mut req = self.http.get(&url); if !query.is_empty() { req = req.query(query); } let resp = req .send() .await .map_err(|e| Error::Transport(e.to_string()))?; let status = resp.status(); if !status.is_success() { let body = resp.text().await.unwrap_or_default(); return Err(map_status(PROVIDER_NAME, status.as_u16(), &body)); } resp.json::<T>() .await .map_err(|e| Error::Transport(e.to_string())) } pub(crate) async fn post_json<B: serde::Serialize, T: DeserializeOwned>( &self, path: &str, body: &B, ) -> Result<T> { let url = format!("{}{}", self.base_url, path); let resp = self .http .post(&url) .json(body) .send() .await .map_err(|e| Error::Transport(e.to_string()))?; let status = resp.status(); if !status.is_success() { let body = resp.text().await.unwrap_or_default(); return Err(map_status(PROVIDER_NAME, status.as_u16(), &body)); } resp.json::<T>() .await .map_err(|e| Error::Transport(e.to_string())) } /// Issues, в которых вовлечён текущий пользователь (`GET /me/issues`). pub async fn list_my_issues(&self) -> Result<Page<Issue>> { let api: ApiIssuesResponse = self.get_json("/me/issues", &[]).await?; let mut page = Page::new(api.issues.into_iter().map(|i| i.into_issue()).collect()); if !api.next_page_token.is_empty() { page = page.with_next(api.next_page_token); } Ok(page) } } fn map_status(provider: &str, status: u16, body: &str) -> Error { match status { 401 => Error::Unauthorized(body.to_string()), 403 => Error::Forbidden(body.to_string()), 404 => Error::NotFound(body.to_string()), 429 => Error::RateLimited { retry_after_secs: 60, }, _ => Error::Provider { provider: provider.to_string(), status, message: body.to_string(), }, } } pub(crate) fn parse_scope(s: &str) -> Result<(&str, &str)> { let mut parts = s.splitn(2, '/'); let owner = parts .next() .ok_or_else(|| Error::InvalidInput(format!("scope `{s}` is not org/repo")))?; let repo = parts .next() .ok_or_else(|| Error::InvalidInput(format!("scope `{s}` is not org/repo")))?; if owner.is_empty() || repo.is_empty() { return Err(Error::InvalidInput(format!("scope `{s}` is not org/repo"))); } Ok((owner, repo)) } impl Provider for SourceCraftClient { fn provider_name(&self) -> &str { PROVIDER_NAME } } #[async_trait] impl RepoProvider for SourceCraftClient { async fn list_repos(&self, filter: RepoFilter) -> Result<Page<Repo>> { let owner = filter.owner.as_deref().ok_or_else(|| { Error::InvalidInput( "filter.owner (organization slug) is required for SourceCraft list_repos".into(), ) })?; let path = format!("/orgs/{owner}/repos"); let mut query: Vec<(String, String)> = Vec::new(); if let Some(limit) = filter.pagination.limit { query.push(("page_size".to_string(), limit.to_string())); } if let Some(cursor) = filter.pagination.cursor.as_deref() { query.push(("page_token".to_string(), cursor.to_string())); } let api: ApiRepositoriesResponse = self.get_json(&path, &query).await?; let mut page = Page::new(api.repositories.into_iter().map(Into::into).collect()); if !api.next_page_token.is_empty() { page = page.with_next(api.next_page_token); } Ok(page) } async fn get_repo(&self, owner: &str, repo: &str) -> Result<Repo> { let path = format!("/repos/{owner}/{repo}"); let api: ApiRepository = self.get_json(&path, &[]).await?; Ok(api.into()) } async fn list_branches(&self, owner: &str, repo: &str) -> Result<Page<Branch>> { let path = format!("/repos/{owner}/{repo}/branches"); let api: ApiBranchesResponse = self.get_json(&path, &[]).await?; let mut page = Page::new(api.branches.into_iter().map(Into::into).collect()); if !api.next_page_token.is_empty() { page = page.with_next(api.next_page_token); } Ok(page) } async fn list_commits( &self, _owner: &str, _repo: &str, _opts: CommitsOpts, ) -> Result<Page<Commit>> { Err(Error::Provider { provider: PROVIDER_NAME.into(), status: 501, message: "list_commits is not exposed via SourceCraft REST API; use git on a clone" .into(), }) } } #[async_trait] impl IssueProvider for SourceCraftClient { async fn list_issues(&self, filter: IssueFilter) -> Result<Page<Issue>> { let scope = filter.scope.as_deref().ok_or_else(|| { Error::InvalidInput("scope (org/repo) is required for SourceCraft list_issues".into()) })?; let (owner, repo) = parse_scope(scope)?; let path = format!("/repos/{owner}/{repo}/issues"); let mut query: Vec<(String, String)> = Vec::new(); if let Some(limit) = filter.pagination.limit { query.push(("page_size".to_string(), limit.to_string())); } if let Some(cursor) = filter.pagination.cursor.as_deref() { query.push(("page_token".to_string(), cursor.to_string())); } let scope_owned = scope.to_string(); let api: ApiIssuesResponse = self.get_json(&path, &query).await?; let mut page = Page::new( api.issues .into_iter() .map(|i| i.into_issue_with_scope(Some(&scope_owned))) .collect(), ); if !api.next_page_token.is_empty() { page = page.with_next(api.next_page_token); } Ok(page) } async fn get_issue(&self, key: &IssueKey) -> Result<Issue> { let (owner, repo) = parse_scope(&key.scope)?; let path = format!("/repos/{owner}/{repo}/issues/{}", key.id); let scope_owned = key.scope.clone(); let api: ApiIssue = self.get_json(&path, &[]).await?; Ok(api.into_issue_with_scope(Some(&scope_owned))) } async fn create_issue(&self, input: CreateIssueInput) -> Result<Issue> { let (owner, repo) = parse_scope(&input.scope)?; let path = format!("/repos/{owner}/{repo}/issues"); let body = serde_json::json!({ "title": input.title, "description": input.body.unwrap_or_default(), }); let scope_owned = input.scope.clone(); let api: ApiIssue = self.post_json(&path, &body).await?; Ok(api.into_issue_with_scope(Some(&scope_owned))) } async fn list_comments(&self, key: &IssueKey) -> Result<Page<Comment>> { let (owner, repo) = parse_scope(&key.scope)?; let path = format!("/repos/{owner}/{repo}/issues/{}/comments", key.id); let api: ApiCommentsResponse = self.get_json(&path, &[]).await?; let mut page = Page::new(api.comments.into_iter().map(Into::into).collect()); if !api.next_page_token.is_empty() { page = page.with_next(api.next_page_token); } Ok(page) } async fn add_comment(&self, key: &IssueKey, body: &str) -> Result<Comment> { let (owner, repo) = parse_scope(&key.scope)?; let path = format!("/repos/{owner}/{repo}/issues/{}/comments", key.id); let payload = serde_json::json!({ "body": body }); let api: ApiCommentRecord = self.post_json(&path, &payload).await?; Ok(api.into()) } } #[async_trait] impl MergeRequestProvider for SourceCraftClient { async fn list_merge_requests(&self, filter: MrFilter) -> Result<Page<MergeRequest>> { let scope = filter.scope.as_deref().ok_or_else(|| { Error::InvalidInput( "scope (org/repo) is required for SourceCraft list_merge_requests".into(), ) })?; let (owner, repo) = parse_scope(scope)?; let path = format!("/repos/{owner}/{repo}/pulls"); let mut query: Vec<(String, String)> = Vec::new(); if let Some(limit) = filter.pagination.limit { query.push(("page_size".to_string(), limit.to_string())); } let scope_owned = scope.to_string(); let api: ApiPullRequestsResponse = self.get_json(&path, &query).await?; let mut page = Page::new( api.pull_requests .into_iter() .map(|p| p.into_merge_request(Some(&scope_owned))) .collect(), ); if !api.next_page_token.is_empty() { page = page.with_next(api.next_page_token); } Ok(page) } async fn get_merge_request( &self, owner: &str, repo: &str, number: u64, ) -> Result<MergeRequest> { let path = format!("/repos/{owner}/{repo}/pulls/{number}"); let scope = format!("{owner}/{repo}"); let api: ApiPullRequest = self.get_json(&path, &[]).await?; Ok(api.into_merge_request(Some(&scope))) } } #[cfg(test)] mod tests { use super::*; use httpmock::prelude::*; #[tokio::test] async fn provider_name_is_sourcecraft() { let c = SourceCraftClient::with_base_url(SecretString::from("t"), "http://127.0.0.1:1") .unwrap(); assert_eq!(c.provider_name(), "sourcecraft"); } #[tokio::test] async fn list_my_issues_parses_response() { let server = MockServer::start_async().await; let _m = server .mock_async(|when, then| { when.method(GET) .path("/me/issues") .header("authorization", "Bearer pv1_t"); then.status(200).json_body(serde_json::json!({ "issues": [ { "id": "uuid-1", "slug": "1", "title": "Bug", "description": "broken", "status": { "id": "1", "slug": "open", "name": "Open", "status_type": "initial" }, "author": { "id": "u1", "slug": "octo" }, "labels": [{"slug": "bug"}], "priority": "normal", "repository": { "id": "r1", "slug": "alpha", "name": "alpha", "organization": { "id": "o1", "slug": "octo" } }, "created_at": "2026-05-01T00:00:00Z" } ], "next_page_token": "" })); }) .await; let client = SourceCraftClient::with_base_url(SecretString::from("pv1_t"), server.base_url()) .unwrap(); let page = client.list_my_issues().await.unwrap(); assert_eq!(page.items.len(), 1); assert_eq!(page.items[0].title, "Bug"); assert_eq!(page.items[0].state, codenik_core::IssueState::Open); assert_eq!(page.items[0].key.scope, "octo/alpha"); assert_eq!(page.items[0].key.id, "1"); assert_eq!(page.items[0].labels, vec!["bug".to_string()]); assert!(page.next.is_none()); } #[tokio::test] async fn unauthorized_response_maps_to_unauthorized_error() { let server = MockServer::start_async().await; let _m = server .mock_async(|when, then| { when.method(GET).path("/me/issues"); then.status(401).body("bad token"); }) .await; let client = SourceCraftClient::with_base_url(SecretString::from("bad"), server.base_url()).unwrap(); let err = client.list_my_issues().await.unwrap_err(); assert!(matches!(err, Error::Unauthorized(_))); } #[tokio::test] async fn list_repos_parses_response() { let server = MockServer::start_async().await; let _m = server .mock_async(|when, then| { when.method(GET).path("/orgs/octo/repos"); then.status(200).json_body(serde_json::json!({ "repositories": [ { "id": "r1", "slug": "alpha", "name": "alpha", "organization": {"id": "o1", "slug": "octo"}, "default_branch": "main", "visibility": "public", "clone_url": { "https": "https://example.org/octo/alpha.git", "ssh": "ssh://git@example.org/octo/alpha.git" }, "web_url": "https://example.org/octo/alpha" } ], "next_page_token": "next-1" })); }) .await; let client = SourceCraftClient::with_base_url(SecretString::from("t"), server.base_url()).unwrap(); let page = client .list_repos(RepoFilter { owner: Some("octo".into()), ..Default::default() }) .await .unwrap(); assert_eq!(page.items.len(), 1); assert_eq!(page.items[0].full_name, "octo/alpha"); assert!(!page.items[0].is_private); assert_eq!(page.next.as_deref(), Some("next-1")); } #[tokio::test] async fn list_repos_without_owner_returns_invalid_input() { let c = SourceCraftClient::with_base_url(SecretString::from("t"), "http://127.0.0.1:1") .unwrap(); let err = c.list_repos(RepoFilter::default()).await.unwrap_err(); assert!(matches!(err, Error::InvalidInput(_))); } #[tokio::test] async fn get_repo_parses_response() { let server = MockServer::start_async().await; let _m = server .mock_async(|when, then| { when.method(GET).path("/repos/octo/alpha"); then.status(200).json_body(serde_json::json!({ "id": "r1", "slug": "alpha", "name": "alpha", "organization": {"id": "o1", "slug": "octo"}, "default_branch": "main", "visibility": "internal" })); }) .await; let client = SourceCraftClient::with_base_url(SecretString::from("t"), server.base_url()).unwrap(); let r = client.get_repo("octo", "alpha").await.unwrap(); assert_eq!(r.full_name, "octo/alpha"); assert!(r.is_private); } #[tokio::test] async fn list_branches_parses_response() { let server = MockServer::start_async().await; let _m = server .mock_async(|when, then| { when.method(GET).path("/repos/octo/alpha/branches"); then.status(200).json_body(serde_json::json!({ "branches": [ { "name": "main", "commit": { "hash": "abc" } }, { "name": "dev", "commit": { "hash": "def" } } ], "next_page_token": "" })); }) .await; let client = SourceCraftClient::with_base_url(SecretString::from("t"), server.base_url()).unwrap(); let page = client.list_branches("octo", "alpha").await.unwrap(); assert_eq!(page.items.len(), 2); assert_eq!(page.items[0].commit_sha, "abc"); } #[tokio::test] async fn list_commits_returns_provider_error() { let c = SourceCraftClient::with_base_url(SecretString::from("t"), "http://127.0.0.1:1") .unwrap(); let err = c .list_commits("a", "b", CommitsOpts::default()) .await .unwrap_err(); match err { Error::Provider { status, .. } => assert_eq!(status, 501), _ => panic!("expected Provider error"), } } #[tokio::test] async fn list_issues_parses_response_with_scope() { let server = MockServer::start_async().await; let _m = server .mock_async(|when, then| { when.method(GET).path("/repos/octo/alpha/issues"); then.status(200).json_body(serde_json::json!({ "issues": [{ "id": "uuid-1", "slug": "1", "title": "Bug", "status": {"id": "1", "slug": "open"}, "author": {"id": "u1", "slug": "octo"} }], "next_page_token": "" })); }) .await; let client = SourceCraftClient::with_base_url(SecretString::from("t"), server.base_url()).unwrap(); let page = client .list_issues(IssueFilter { scope: Some("octo/alpha".into()), ..Default::default() }) .await .unwrap(); assert_eq!(page.items.len(), 1); assert_eq!(page.items[0].key.scope, "octo/alpha"); } #[tokio::test] async fn add_comment_posts_body() { let server = MockServer::start_async().await; let _m = server .mock_async(|when, then| { when.method(POST) .path("/repos/octo/alpha/issues/7/comments") .json_body(serde_json::json!({"body": "hi"})); then.status(201).json_body(serde_json::json!({ "id": "c1", "body": "hi" })); }) .await; let client = SourceCraftClient::with_base_url(SecretString::from("t"), server.base_url()).unwrap(); let key = IssueKey::new("octo/alpha", "7"); let c = client.add_comment(&key, "hi").await.unwrap(); assert_eq!(c.id, "c1"); assert_eq!(c.body, "hi"); } #[tokio::test] async fn list_pulls_parses_response() { let server = MockServer::start_async().await; let _m = server .mock_async(|when, then| { when.method(GET).path("/repos/octo/alpha/pulls"); then.status(200).json_body(serde_json::json!({ "pull_requests": [ { "id": "p1", "slug": "5", "title": "Add feature", "status": {"id": "1", "slug": "open"}, "author": {"id": "u1", "slug": "octo"}, "source_branch": "feat/x", "target_branch": "main" } ], "next_page_token": "" })); }) .await; let client = SourceCraftClient::with_base_url(SecretString::from("t"), server.base_url()).unwrap(); let page = client .list_merge_requests(MrFilter { scope: Some("octo/alpha".into()), ..Default::default() }) .await .unwrap(); assert_eq!(page.items.len(), 1); assert_eq!(page.items[0].key.scope, "octo/alpha"); assert_eq!(page.items[0].source_branch, "feat/x"); assert_eq!(page.items[0].state, codenik_core::MergeRequestState::Open); } /// Smoke: проверяем доступность реальных эндпоинтов с PAT. #[tokio::test] #[ignore = "smoke test against real SourceCraft API"] async fn smoke_list_my_issues_against_real_api() { let Ok(token) = std::env::var("SOURCECRAFT_TOKEN") else { return; }; if token.is_empty() { return; } let client = SourceCraftClient::new(SecretString::from(token)).unwrap(); let page = client .list_my_issues() .await .expect("list_my_issues failed"); eprintln!("smoke: /me/issues returned {} item(s)", page.items.len()); } #[tokio::test] #[ignore = "smoke test against real SourceCraft API"] async fn smoke_list_repos_against_real_api() { let Ok(token) = std::env::var("SOURCECRAFT_TOKEN") else { return; }; if token.is_empty() { return; } let client = SourceCraftClient::new(SecretString::from(token)).unwrap(); let owner = std::env::var("SOURCECRAFT_SMOKE_ORG").unwrap_or_else(|_| "codenik".into()); let page = client .list_repos(RepoFilter { owner: Some(owner), ..Default::default() }) .await .expect("list_repos failed"); eprintln!( "smoke: list_repos returned {} repo(s)", page.items.len() ); } }