32 Commits
Author SHA1 Message Date
qpismont a29051b0e4 Fix clippy errors
ci/woodpecker/push/tests Pipeline was successful
2026-07-31 20:39:50 +00:00
qpismont 8c53bc0e20 re fix fmt lol
ci/woodpecker/push/tests Pipeline failed
2026-07-31 20:33:47 +00:00
qpismont f0e64e0c1d Fix fmt
ci/woodpecker/push/tests Pipeline failed
2026-07-31 20:32:59 +00:00
qpismont 6a21c7d6c3 Renforce woodpecker tests
ci/woodpecker/push/tests Pipeline failed
2026-07-31 20:25:37 +00:00
qpismont b3a0cb63e9 Update woodpecker rust job (1.96 => 1.97)
ci/woodpecker/push/tests Pipeline was successful
2026-07-31 20:21:05 +00:00
qpismont 5b9d870b46 Dockerfile field must be present
ci/woodpecker/push/tests Pipeline was successful
2026-07-31 20:18:29 +00:00
qpismont 15f619ccf7 Move to multi crates project
ci/woodpecker/push/tests Pipeline was successful
Starting impl devcontainer spec
2026-07-31 20:06:37 +00:00
qpismont d711153553 Merge pull request 'Observability' (#6) from 1.1 into main
ci/woodpecker/push/tests Pipeline was successful
ci/woodpecker/tag/release Pipeline was successful
Reviewed-on: #6
2026-07-26 21:55:41 +02:00
qpismont b6a299ac18 bump version
ci/woodpecker/push/tests Pipeline was successful
2026-07-26 18:41:22 +00:00
qpismont c426b9b513 update readme
ci/woodpecker/push/tests Pipeline was successful
2026-07-26 11:56:58 +00:00
qpismont e1cb5d7d96 fix usd metric
ci/woodpecker/push/tests Pipeline was successful
2026-07-26 11:41:58 +00:00
qpismont 6ffd88927c add prometheus metrics
ci/woodpecker/push/tests Pipeline was successful
2026-07-26 11:02:34 +00:00
qpismont 7252bf7673 fix tests
ci/woodpecker/push/tests Pipeline was successful
2026-06-30 20:48:46 +00:00
qpismont 743b6b33c9 Switch to vscode + fetch bot_name with token
ci/woodpecker/push/tests Pipeline failed
2026-06-30 20:44:25 +00:00
qpismont 7f24d7657c Merge pull request 'Fix missing env var error' (#5) from 1.0.1 into main
ci/woodpecker/tag/release Pipeline was successful
Reviewed-on: #5
2026-06-16 22:05:24 +02:00
qpismont a613fdb99e Another trace clear 2026-06-16 18:53:33 +00:00
qpismont 975581093a Add missing sentry error backtrace
Clear traces spam
2026-06-16 18:49:35 +00:00
qpismont 00d46ce968 Fix AI review
Webhook action check before user check
2026-06-12 22:01:49 +00:00
qpismont 3f6c5b5559 Fix missing env var error
ci/woodpecker/push/tests Pipeline was successful
2026-06-12 21:27:12 +00:00
qpismont d4666fb36e Merge pull request 'prepare first release with graceful shutdown + containerfile + push to' (#4) from 1.0 into main
ci/woodpecker/push/tests Pipeline was successful
ci/woodpecker/tag/release Pipeline was successful
Reviewed-on: #4
2026-06-12 22:38:25 +02:00
qpismont cf59455d4a Fix release job
ci/woodpecker/push/tests Pipeline was successful
2026-06-12 20:15:20 +00:00
qpismont c7387a3b28 Fix tests job
ci/woodpecker/push/tests Pipeline was successful
2026-06-12 19:57:28 +00:00
qpismont 433021d607 Add webhook already handled check
ci/woodpecker/push/tests Pipeline failed
ci/woodpecker/pr/tests Pipeline failed
Fix all tests
Add woodpecker ci (tests + release)
2026-06-12 19:56:32 +00:00
qpismont 3c32cd20b6 Readme :) 2026-06-10 20:01:21 +00:00
qpismont a30d7c5d90 prepare first release with graceful shutdown + containerfile + push to
hub script
2026-06-10 19:23:17 +00:00
qpismont 9175f9b3a2 Merge pull request 'Starting impl Sentry and tracing' (#3) from tracing into main
Reviewed-on: #3
2026-06-10 20:27:26 +02:00
qpismont 71ebfdd276 tasks.join_next => join_all() 2026-06-10 18:26:55 +00:00
qpismont 3d751ae6c6 drain completed tasks and log webhook queue stats 2026-06-10 18:19:18 +00:00
qpismont 6599c20c30 verify_signature before adding body to sentry event 2026-06-10 18:08:35 +00:00
qpismont a2d898c07d Fix sentry http request info
Add async bot running with semaphore
2026-06-10 17:50:24 +00:00
qpismont efb35d5a8a continue tracing impl 2026-06-09 21:17:03 +00:00
qpismont 39c2afa0a7 Starting impl Sentry and tracing 2026-06-09 20:58:38 +00:00
30 changed files with 2740 additions and 366 deletions
+7 -6
View File
@@ -1,4 +1,4 @@
FROM debian:trixie
FROM rust:1.97-trixie
ARG USERNAME=dev
ARG USER_UID=1000
@@ -11,17 +11,18 @@ RUN apt-get update && apt-get install -y --no-install-recommends \
curl \
git \
build-essential \
libssl-dev \
cmake \
pkg-config \
clang \
&& rm -rf /var/lib/apt/lists/*
RUN groupadd --gid ${USER_GID:-1000} $USERNAME \
&& useradd --uid ${USER_UID:-1000} --gid ${USER_GID:-1000} -m $USERNAME
&& useradd --uid ${USER_UID:-1000} --gid ${USER_GID:-1000} -m $USERNAME \
&& rustup component add clippy \
&& rustup component add rustfmt
USER $USERNAME
WORKDIR /home/$USERNAME
ENV PATH="/home/${USERNAME}/.cargo/bin:${PATH}"
RUN curl --proto '=https' --tlsv1.2 -sSf https://sh.rustup.rs | sh -s -- -y --default-toolchain stable
+6 -1
View File
@@ -12,8 +12,13 @@
"containerEnv": {
"SHELL": "/bin/bash"
},
"customizations": {
"vscode": {
"extensions": ["fill-labs.dependi", "rust-lang.rust-analyzer", "tamasfe.even-better-toml"]
}
},
"workspaceMount": "source=${localWorkspaceFolder},target=/workspaces/herald,type=bind",
"workspaceFolder": "/workspaces/herald",
"runArgs": ["--userns=keep-id", "--security-opt", "label=disable"],
"runArgs": ["--userns=keep-id", "--security-opt", "label=disable"],
"appPort": [3000]
}
+22
View File
@@ -0,0 +1,22 @@
HTTP_PORT=3000
BOT_NAME=Herald
WEBHOOK_SIG_HEADER_SECRET=
OPEN_ROUTER_API_KEY=
OPEN_ROUTER_MODEL=deepseek/deepseek-v4-flash
OPEN_ROUTER_TIMEOUT=120
BOT_MAX_CONCURRENT=5
GITEA_URL=https://gitea.example.com
GITEA_TOKEN=
GITEA_TIMEOUT=30
# Optional
SENTRY_DSN=
RUST_LOG=info
RUST_BACKTRACE=1
METRICS_BIND_ADDR=
+19
View File
@@ -0,0 +1,19 @@
when:
event:
- tag
steps:
- name: release-docker
image: quay.io/buildah/stable
privileged: true
volumes:
- /data/woodpecker-builds:/data
commands:
- echo $DOCKER_PASSWORD | buildah login docker.io -u $DOCKER_USERNAME --password-stdin
- chmod +x scripts/build.sh
- bash scripts/build.sh
environment:
DOCKER_USERNAME:
from_secret: docker_username
DOCKER_PASSWORD:
from_secret: docker_password
+29
View File
@@ -0,0 +1,29 @@
when:
event:
- push
steps:
- name: fmt
image: rust:1.97
commands:
- rustup component add rustfmt
- cargo fmt --all -- --check
- name: clippy
image: rust:1.97
commands:
- rustup component add clippy
- cargo clippy --workspace --all-targets --all-features -- -D warnings
- name: test
image: rust:1.97
commands:
- cargo test --workspace --all-targets
- name: container-build
image: quay.io/buildah/stable
privileged: true
volumes:
- /data/woodpecker-builds:/data
commands:
- buildah bud -f Containerfile -t herald-ci .
+22
View File
@@ -0,0 +1,22 @@
{
"languages": {
"Rust": {
"format_on_save": "on",
"formatter": "language_server"
}
},
"lsp": {
"rust-analyzer": {
"initialization_options": {
"check": {
"command": "clippy",
"extraArgs": [
"--",
"-D",
"warnings"
]
}
}
}
}
}
Generated
+1555 -83
View File
File diff suppressed because it is too large Load Diff
+23 -12
View File
@@ -1,23 +1,34 @@
[package]
name = "herald"
version = "0.1.0"
edition = "2024"
[workspace]
members = [
"crates/herald-server",
"crates/devcontainer-rs",
]
resolver = "3"
[dependencies]
[workspace.dependencies]
reqwest = { version = "0.12", default-features = false, features = ["json", "rustls-tls"] }
tokio = { version = "1.52", features = ["full"] }
tokio = { version = "1.53", features = ["full"] }
tokio-stream = "0.1"
tokio-util = "0.7"
futures-util = "0.3"
serde_json = "1.0"
serde = { version = "1.0", features = ["derive"] }
openrouter-rs = "0.10"
sentry = { version = "0.48", features = ["tower-axum-matched-path"] }
sentry-anyhow = { version = "0.48", features = ["backtrace"] }
openrouter-rs = "0.12"
dotenvy = "0.15"
tower = "0.5"
tower-http = { version = "0.6", features = ["trace"] }
tracing = "0.1"
tracing-subscriber = { version = "0.3", features = ["env-filter"] }
axum = "0.8"
anyhow = "1.0"
anyhow = { version = "1.0", features = ["backtrace"] }
thiserror = "2.0"
hmac = "0.13"
sha2 = "0.11"
ring = "0.17"
hex = "0.4"
subtle = "2.6"
bytes = "1.11"
bytes = "1.1"
metrics = "0.24"
metrics-exporter-prometheus = { version = "0.18", default-features = false, features = ["http-listener"] }
[profile.release]
debug = 1
+15
View File
@@ -0,0 +1,15 @@
FROM rust:1.97-trixie as builder
WORKDIR /app
COPY Cargo.toml Cargo.lock ./
COPY crates/ crates/
RUN cargo build --release --package herald-server
FROM debian:trixie-slim
WORKDIR /app
COPY --from=builder /app/target/release/herald-server .
CMD [ "./herald-server" ]
+92
View File
@@ -0,0 +1,92 @@
# Herald
Herald is a Gitea bot that performs automated AI-powered code reviews on pull requests using [OpenRouter](https://openrouter.ai/).
## Features
- Listens for Gitea webhook events and triggers code reviews on pull request comments
- Streams reviews back to Gitea as PR comments
- Concurrent review processing with configurable parallelism
- Graceful shutdown — in-progress reviews finish before the process exits
- Error tracking via Sentry
- Prometheus metrics endpoint for monitoring
- Tiny memory footprint (~4MB) thanks to Rust
## Installation
**Requirements:** Rust toolchain ([rustup.rs](https://rustup.rs))
```sh
cargo build --release
./target/release/herald
```
Herald reads its configuration from environment variables (a `.env` file is supported):
| Variable | Description |
|---|---|
| `HTTP_PORT` | Port to listen on |
| `BOT_NAME` | The bot's Gitea username (used to detect mentions) |
| `WEBHOOK_SIG_HEADER_SECRET` | Gitea webhook secret for signature verification |
| `OPEN_ROUTER_API_KEY` | OpenRouter API key |
| `OPEN_ROUTER_MODEL` | Model to use (e.g. `deepseek/deepseek-v4-flash`) |
| `OPEN_ROUTER_TIMEOUT` | OpenRouter request timeout in seconds |
| `BOT_MAX_CONCURRENT` | Maximum number of concurrent reviews |
| `GITEA_URL` | Base URL of your Gitea instance |
| `GITEA_TOKEN` | Gitea API token |
| `GITEA_TIMEOUT` | Gitea API request timeout in seconds |
| `METRICS_BIND_ADDR` | *(optional)* Bind address for the Prometheus metrics endpoint (e.g. `0.0.0.0:9100`). If unset, the metrics exporter is disabled. |
| `SENTRY_DSN` | *(optional)* Sentry DSN for error tracking |
| `RUST_LOG` | *(optional)* Log level, defaults to `info` |
## Development
The easiest way to get started is with the provided [Dev Container](https://containers.dev/) (VS Code or Zed with the dev container extension).
Open the project and reopen it in the container — the Rust toolchain is pre-installed.
**Without Dev Container**, you just need a Rust toolchain:
```sh
rustup toolchain install stable
cargo run
```
Copy `.env.example` to `.env` and fill in your values before running.
## Metrics
Herald optionally exposes a Prometheus metrics endpoint, useful for scraping with an [OpenTelemetry Collector](https://opentelemetry.io/docs/collector/) or any Prometheus-compatible scraper.
Set `METRICS_BIND_ADDR` (e.g. `0.0.0.0:9100`) to enable it. The metrics are then available at `http://<host>:9100/metrics`.
### Exposed metrics
| Metric | Type | Description |
|---|---|---|
| `herald_webhooks_received_total` | counter | Total webhooks received (label: `event_type`) |
| `herald_webhooks_duplicate_total` | counter | Webhooks rejected as duplicates (label: `event_type`) |
| `herald_webhooks_channel_full_total` | counter | Webhooks dropped because the bot channel was full (label: `event_type`) |
| `herald_bot_tasks_active` | gauge | Bot tasks currently in progress |
| `herald_bot_tasks_completed_total` | counter | Bot tasks completed successfully (label: `event_type`) |
| `herald_bot_tasks_failed_total` | counter | Bot tasks that failed (label: `event_type`) |
| `herald_openrouter_cost_cents_total` | counter | Total OpenRouter cost in cents (divide by 100 for USD) |
### OTel collector example
```yaml
receivers:
prometheus:
config:
scrape_configs:
- job_name: herald
scrape_interval: 15s
static_configs:
- targets: ["herald:9100"]
service:
pipelines:
metrics:
receivers: [prometheus]
exporters: [otlp]
```
View File
+14
View File
@@ -0,0 +1,14 @@
[package]
name = "devcontainer-rs"
version = "0.1.0"
edition = "2024"
[dependencies]
tokio = { workspace = true }
serde = { workspace = true }
serde_json = { workspace = true }
anyhow = { workspace = true }
thiserror = { workspace = true }
[dev-dependencies]
tempfile = "3"
+152
View File
@@ -0,0 +1,152 @@
use std::{
collections::HashMap,
path::{Path, PathBuf},
};
use serde::Deserialize;
#[derive(Debug, Deserialize)]
pub struct DevContainerBuildSchema {
#[serde(default)]
pub dockerfile: String,
#[serde(default)]
pub args: HashMap<String, String>,
}
#[derive(Debug, Deserialize)]
pub struct DevContainerSchema {
#[serde(default)]
pub name: Option<String>,
pub build: DevContainerBuildSchema,
#[serde(rename = "workspaceFolder", default)]
pub workspace_folder: Option<String>,
#[serde(rename = "containerEnv", default)]
pub container_env: HashMap<String, String>,
#[serde(rename = "postCreateCommand", default)]
pub post_create_command: Option<String>,
#[serde(rename = "postStartCommand", default)]
pub post_start_command: Option<String>,
}
#[derive(Debug)]
pub struct DevContainer {
pub container_file_path: PathBuf,
pub name: Option<String>,
pub build_args: HashMap<String, String>,
pub container_env: HashMap<String, String>,
pub workspace_folder: Option<String>,
pub post_create_command: Option<String>,
pub post_start_command: Option<String>,
}
#[derive(Debug, thiserror::Error)]
pub enum ParseError {
#[error("failed to read devcontainer file `{path}`: {source}")]
Read {
path: PathBuf,
source: std::io::Error,
},
#[error("invalid devcontainer JSON in `{path}`: {source}")]
Json {
path: PathBuf,
source: serde_json::Error,
},
#[error("container file `{0}` does not exist or is not a regular file")]
ContainerFileNotFound(PathBuf),
#[error("the devcontainer file path has no parent directory: `{0}`")]
InvalidDevContainerPath(PathBuf),
}
impl TryFrom<(DevContainerSchema, PathBuf)> for DevContainer {
type Error = ParseError;
fn try_from(
(schema, devcontainer_path): (DevContainerSchema, PathBuf),
) -> Result<Self, Self::Error> {
let base_dir = devcontainer_path
.parent()
.ok_or_else(|| ParseError::InvalidDevContainerPath(devcontainer_path.clone()))?;
let container_file_path = base_dir.join(schema.build.dockerfile);
if !container_file_path.is_file() {
return Err(ParseError::ContainerFileNotFound(container_file_path));
}
Ok(Self {
container_file_path,
name: schema.name,
build_args: schema.build.args,
container_env: schema.container_env,
workspace_folder: schema.workspace_folder,
post_create_command: schema.post_create_command,
post_start_command: schema.post_start_command,
})
}
}
pub async fn parse(path: impl AsRef<Path>) -> Result<DevContainer, ParseError> {
let path = path.as_ref().to_path_buf();
let contents = tokio::fs::read_to_string(&path)
.await
.map_err(|source| ParseError::Read {
path: path.clone(),
source,
})?;
let schema = serde_json::from_str::<DevContainerSchema>(&contents).map_err(|source| {
ParseError::Json {
path: path.clone(),
source,
}
})?;
DevContainer::try_from((schema, path))
}
#[cfg(test)]
mod tests {
use super::*;
use std::fs;
#[tokio::test]
async fn parses_devcontainer_file() {
let dir = tempfile::tempdir().unwrap();
let devcontainer_path = dir.path().join("devcontainer.json");
let dockerfile_path = dir.path().join("Dockerfile");
fs::write(&dockerfile_path, "FROM alpine\n").unwrap();
fs::write(
&devcontainer_path,
r#"{
"name": "test",
"build": {
"dockerfile": "Dockerfile",
"args": {
"VERSION": "1"
}
},
"workspaceFolder": "/workspace",
"containerEnv": {
"RUST_LOG": "debug"
}
}"#,
)
.unwrap();
let config = parse(&devcontainer_path).await.unwrap();
assert_eq!(config.name.as_deref(), Some("test"));
assert_eq!(config.container_file_path, dockerfile_path);
assert_eq!(config.build_args.get("VERSION").unwrap(), "1");
assert_eq!(config.container_env.get("RUST_LOG").unwrap(), "debug");
assert_eq!(config.workspace_folder.as_deref(), Some("/workspace"));
}
}
+30
View File
@@ -0,0 +1,30 @@
[package]
name = "herald-server"
version = "1.2.0"
edition = "2024"
[dependencies]
reqwest = { workspace = true }
tokio = { workspace = true }
tokio-stream = { workspace = true }
tokio-util = { workspace = true }
futures-util = { workspace = true }
serde_json = { workspace = true }
serde = { workspace = true }
sentry = { workspace = true }
sentry-anyhow = { workspace = true }
openrouter-rs = { workspace = true }
dotenvy = { workspace = true }
tower = { workspace = true }
tower-http = { workspace = true }
tracing = { workspace = true }
tracing-subscriber = { workspace = true }
axum = { workspace = true }
anyhow = { workspace = true }
thiserror = { workspace = true }
ring = { workspace = true }
hex = { workspace = true }
bytes = { workspace = true }
metrics = { workspace = true }
metrics-exporter-prometheus = { workspace = true }
devcontainer-rs = { path = "../devcontainer-rs" }
+53 -17
View File
@@ -1,46 +1,76 @@
use axum::body::{Bytes, to_bytes};
use axum::extract::{FromRef, FromRequest, State};
use axum::http::Request;
use axum::response::IntoResponse;
use axum::routing::{get, post};
use axum::{Json, Router};
use hmac::{Hmac, KeyInit, Mac};
use reqwest::StatusCode;
use ring::hmac;
use sentry::integrations::tower::{NewSentryLayer, SentryHttpLayer};
use serde_json::Value;
use sha2::Sha256;
use subtle::ConstantTimeEq;
use tower::ServiceBuilder;
use tower_http::trace::TraceLayer;
use tracing::{info, instrument};
use tokio_util::sync::CancellationToken;
use crate::consts::{GITEA_EVENT_TYPE_HEADER_NAME, GITEA_SIG_HEADER_NAME, MAX_WEBHOOK_BODY_SIZE};
use crate::errors::AppError;
use crate::gitea::WebhookType;
use crate::metrics;
use crate::state::AppState;
pub async fn start(app_state: AppState) -> anyhow::Result<()> {
pub async fn start(app_state: AppState, shutdown: CancellationToken) -> anyhow::Result<()> {
let http_port = app_state.config.http_port;
let app = Router::new()
.route("/", get(root))
.route("/webhook", post(webhook))
.layer(
ServiceBuilder::new()
.layer(NewSentryLayer::<Request<_>>::new_from_top())
.layer(SentryHttpLayer::new())
.layer(TraceLayer::new_for_http()),
)
.with_state(app_state);
let listener = tokio::net::TcpListener::bind(format!("0.0.0.0:{}", http_port)).await?;
info!("Listening API on port {}", http_port);
axum::serve(listener, app)
.with_graceful_shutdown(async move { shutdown.cancelled().await })
.await
.map_err(anyhow::Error::from)
.inspect(|_| info!("API shutting down complete"))?;
Ok(())
}
async fn root() -> &'static str {
"Hi, i'm Herald :)"
}
#[instrument(skip(app_state), fields(webhook_type), err)]
async fn webhook(
State(app_state): State<AppState>,
WebhookExtract(wb): WebhookExtract,
) -> Result<impl IntoResponse, AppError> {
app_state
.bot_tx
.send(wb)
.await
.map_err(anyhow::Error::from)?;
tracing::Span::current().record("webhook_type", tracing::field::debug(&wb));
let event_type = wb.event_type_str();
metrics::webhook_received(event_type);
let event_id = wb.event_id();
if app_state.bot.check_and_mark(event_id).await {
metrics::webhook_duplicate(event_type);
return Err(AppError::AlreadyProcessedErr);
}
if app_state.bot_tx.try_send(wb).is_err() {
app_state.bot.unmark(event_id).await;
metrics::webhook_channel_full(event_type);
return Err(AppError::ChannelFullErr);
}
Ok((StatusCode::CREATED, "Task started"))
}
@@ -54,6 +84,7 @@ where
{
type Rejection = AppError;
#[instrument(skip(req, state), err)]
async fn from_request(req: axum::extract::Request, state: &S) -> Result<Self, Self::Rejection> {
let app_state = AppState::from_ref(state);
let headers = req.headers();
@@ -68,7 +99,17 @@ where
&body_bytes,
)?;
let webhook = parse_webhook(&type_header, &app_state.config.bot_name, &body_bytes)?;
let body_str = String::from_utf8_lossy(&body_bytes).into_owned();
sentry::configure_scope(|scope| {
scope.add_event_processor(move |mut event| {
let mut request = event.request.take().unwrap_or_default();
request.data = Some(body_str.clone());
event.request = Some(request);
Some(event)
});
});
let webhook = parse_webhook(&type_header, &app_state.bot.name(), &body_bytes)?;
Ok(WebhookExtract(webhook))
}
}
@@ -100,12 +141,7 @@ fn parse_webhook(header: &str, bot_name: &str, body_bytes: &[u8]) -> Result<Webh
fn verify_signature(secret_key: &[u8], sig_header: &str, body: &[u8]) -> Result<(), AppError> {
let sig_header_decoded =
hex::decode(sig_header).map_err(|_| AppError::WebHookSigHeaderInvalidErr)?;
let mut mac = Hmac::<Sha256>::new_from_slice(secret_key).map_err(anyhow::Error::from)?;
let key = hmac::Key::new(hmac::HMAC_SHA256, secret_key);
mac.update(body);
let generated_hmac = mac.finalize().into_bytes();
bool::from(generated_hmac.ct_eq(&sig_header_decoded))
.then_some(())
.ok_or(AppError::WebHookSigHeaderInvalidErr)
hmac::verify(&key, body, &sig_header_decoded).map_err(|_| AppError::WebHookSigHeaderInvalidErr)
}
+222
View File
@@ -0,0 +1,222 @@
use crate::{
gitea::{GiteaAPI, WebhookType},
metrics,
open_router::OpenRouterClient,
};
use serde::Deserialize;
use std::{collections::HashSet, sync::Arc};
use tokio::sync::Mutex;
use tokio_util::sync::CancellationToken;
use tracing::{error, info, instrument};
#[derive(Deserialize, Debug)]
pub struct ReviewResult {
pub reviews: Vec<ReviewItem>,
pub comment: String,
pub cost: Option<f64>,
}
#[derive(Deserialize, Debug)]
pub struct ReviewItem {
pub filename: String,
pub line: Option<u64>,
pub message: String,
}
#[derive(Clone)]
pub struct Bot {
bot_name: String,
gitea_api: GiteaAPI,
open_router_client: OpenRouterClient,
http_client: reqwest::Client,
max_concurrent: usize,
open_router_model: String,
actions_handled: Arc<Mutex<HashSet<u64>>>,
}
impl Bot {
pub fn new(
bot_name: String,
gitea_api: GiteaAPI,
open_router_client: OpenRouterClient,
http_client: reqwest::Client,
max_concurrent: usize,
open_router_model: String,
) -> Self {
Self {
bot_name,
gitea_api,
open_router_client,
http_client,
max_concurrent,
open_router_model,
actions_handled: Arc::new(Mutex::new(HashSet::new())),
}
}
pub fn name(&self) -> String {
self.bot_name.clone()
}
pub async fn start(
&self,
mut rx: tokio::sync::mpsc::Receiver<WebhookType>,
shutdown: CancellationToken,
) -> anyhow::Result<()> {
info!("Bot started");
let sem = Arc::new(tokio::sync::Semaphore::new(self.max_concurrent));
let mut tasks = tokio::task::JoinSet::new();
loop {
let wb = tokio::select! {
biased;
_ = shutdown.cancelled() => break,
msg = rx.recv() => match msg {
Some(wb) => wb,
None => break,
},
};
while let Some(res) = tasks.try_join_next() {
if let Err(e) = res {
error!("Task panicked: {e}");
}
}
info!(queued = rx.len(), active = tasks.len(), "Webhook received");
let permit = sem.clone().acquire_owned().await?;
let self_clone = self.clone();
metrics::increment_task_active();
tasks.spawn(async move {
self_clone.exec(wb).await;
drop(permit);
metrics::decrement_task_active();
});
}
tasks.join_all().await;
info!("Bot shutting down complete");
Ok(())
}
#[instrument(skip(self, webhook), fields(repo, pr))]
pub async fn exec(&self, webhook: WebhookType) {
let event_type_str = webhook.event_type_str();
match &webhook {
WebhookType::Review(p) => {
tracing::Span::current().record("repo", &p.repository.full_name);
tracing::Span::current().record("pr", p.pull_request.number);
}
};
let exec_result = match webhook {
WebhookType::Review(review_payload) => crate::bot_actions::review::exec_review(
&self.gitea_api,
&self.open_router_client,
&self.http_client,
&self.open_router_model,
review_payload,
),
}
.await;
match exec_result {
Ok(_) => {
metrics::task_completed(event_type_str);
info!("Task completed");
}
Err(err) => {
metrics::task_failed(event_type_str);
error!(%err, "Task error");
sentry_anyhow::capture_anyhow(&err);
}
}
}
pub async fn check_and_mark(&self, event_id: u64) -> bool {
let mut action_handled_lock = self.actions_handled.lock().await;
!action_handled_lock.insert(event_id)
}
pub async fn unmark(&self, event_id: u64) {
let mut action_handled_lock = self.actions_handled.lock().await;
action_handled_lock.remove(&event_id);
}
}
#[cfg(test)]
mod tests {
use super::*;
fn make_actions_handled() -> Arc<Mutex<HashSet<u64>>> {
Arc::new(Mutex::new(HashSet::new()))
}
async fn check_and_mark(actions_handled: &Arc<Mutex<HashSet<u64>>>, event_id: u64) -> bool {
let mut lock = actions_handled.lock().await;
!lock.insert(event_id)
}
async fn unmark(actions_handled: &Arc<Mutex<HashSet<u64>>>, event_id: u64) {
let mut lock = actions_handled.lock().await;
lock.remove(&event_id);
}
#[tokio::test]
async fn test_check_and_mark_first_call_returns_false() {
let actions_handled = make_actions_handled();
assert!(!check_and_mark(&actions_handled, 1).await);
}
#[tokio::test]
async fn test_check_and_mark_second_call_returns_true() {
let actions_handled = make_actions_handled();
check_and_mark(&actions_handled, 1).await;
assert!(check_and_mark(&actions_handled, 1).await);
}
#[tokio::test]
async fn test_check_and_mark_different_ids_are_independent() {
let actions_handled = make_actions_handled();
check_and_mark(&actions_handled, 1).await;
assert!(!check_and_mark(&actions_handled, 2).await);
}
#[tokio::test]
async fn test_unmark_allows_reprocessing() {
let actions_handled = make_actions_handled();
check_and_mark(&actions_handled, 1).await;
unmark(&actions_handled, 1).await;
assert!(!check_and_mark(&actions_handled, 1).await);
}
#[tokio::test]
async fn test_unmark_nonexistent_id_is_noop() {
let actions_handled = make_actions_handled();
unmark(&actions_handled, 99).await;
assert!(!check_and_mark(&actions_handled, 99).await);
}
#[tokio::test]
async fn test_check_and_mark_concurrent_same_id() {
let actions_handled = make_actions_handled();
let actions_handled2 = Arc::clone(&actions_handled);
let t1 = tokio::spawn(async move { check_and_mark(&actions_handled, 42).await });
let t2 = tokio::spawn(async move { check_and_mark(&actions_handled2, 42).await });
let (r1, r2) = tokio::join!(t1, t2);
let results = [r1.unwrap(), r2.unwrap()];
// exactement un seul des deux doit retourner false (non traité)
assert_eq!(results.iter().filter(|&&r| !r).count(), 1);
// l'autre doit retourner true (déjà traité)
assert_eq!(results.iter().filter(|&&r| r).count(), 1);
}
}
@@ -1,22 +1,31 @@
use futures_util::stream::TryStreamExt;
use tokio::io::AsyncReadExt;
use tokio_util::io::StreamReader;
use tracing::instrument;
use crate::{
bot::ReviewResult,
consts::{BOT_PROCESS_MSG, MAX_DIFF_SIZE, REVIEW_PROMPT},
errors::AppError,
gitea::{GiteaAPI, ReviewPayload},
metrics,
open_router::OpenRouterClient,
};
#[instrument(skip(gitea_api, open_router_client, http_client, review_payload))]
pub async fn exec_review(
gitea_api: &GiteaAPI,
open_router_client: &OpenRouterClient,
http_client: &reqwest::Client,
model: &str,
review_payload: ReviewPayload,
) -> Result<(), AppError> {
) -> anyhow::Result<()> {
tracing::info!(
repo = %review_payload.repository.full_name,
pr = review_payload.pull_request.number,
action = %review_payload.action,
"Starting review"
);
let new_comment = gitea_api
.comment(
&BOT_PROCESS_MSG.replace("{model}", model),
@@ -27,7 +36,7 @@ pub async fn exec_review(
let bot_result: Result<ReviewResult, anyhow::Error> = async {
let git_diff =
download_git_diff(&http_client, &review_payload.pull_request.diff_url).await?;
download_git_diff(http_client, &review_payload.pull_request.diff_url).await?;
let diff_for_llm = format_diff_for_review(&git_diff);
@@ -40,10 +49,16 @@ pub async fn exec_review(
let mut review_result = serde_json::from_str::<ReviewResult>(&chat_result.message)?;
review_result.cost = chat_result.cost;
if let Some(cost) = review_result.cost {
metrics::openrouter_cost_usd(cost);
}
let final_review_markdown = review_result_to_markdown(&review_result);
gitea_api
.post_pull_request_review(
&review_result,
&final_review_markdown,
&review_payload.repository.full_name,
review_payload.pull_request.number,
)
@@ -53,20 +68,22 @@ pub async fn exec_review(
}
.await;
let edit_msg = match bot_result {
Ok(bot_result) => review_result_to_markdown(&bot_result),
Err(e) => format!("Error while reviewing: {}", e),
};
gitea_api
.edit_comment(
&edit_msg,
&review_payload.repository.full_name,
new_comment.id,
)
.await?;
Ok(())
match bot_result {
Ok(_) => {
gitea_api
.delete_comment(&review_payload.repository.full_name, new_comment.id)
.await
}
Err(e) => {
gitea_api
.edit_comment(
&format!("Error while reviewing: {}", e),
&review_payload.repository.full_name,
new_comment.id,
)
.await
}
}
}
fn review_result_to_markdown(review_result: &ReviewResult) -> String {
+82
View File
@@ -0,0 +1,82 @@
use anyhow::anyhow;
#[derive(Clone)]
pub struct EnvConfig {
pub http_port: u16,
pub webhook_secret: String,
pub open_router_api_key: String,
pub open_router_model: String,
pub open_router_timeout: u64,
pub bot_max_concurrent: usize,
pub gitea_url: String,
pub gitea_token: String,
pub gitea_timeout: u64,
pub metrics_bind_addr: Option<String>,
}
pub fn load_config() -> anyhow::Result<EnvConfig> {
let http_port = try_get_env("HTTP_PORT")?.parse()?;
let webhook_secret = try_get_env("WEBHOOK_SIG_HEADER_SECRET")?;
let open_router_api_key = try_get_env("OPEN_ROUTER_API_KEY")?;
let open_router_model = try_get_env("OPEN_ROUTER_MODEL")?;
let open_router_timeout = try_get_env("OPEN_ROUTER_TIMEOUT")?.parse()?;
let bot_max_concurrent = try_get_env("BOT_MAX_CONCURRENT")?.parse()?;
let gitea_url = try_get_env("GITEA_URL")?;
let gitea_token = try_get_env("GITEA_TOKEN")?;
let gitea_timeout = try_get_env("GITEA_TIMEOUT")?.parse()?;
let metrics_bind_addr = std::env::var("METRICS_BIND_ADDR").ok();
Ok(EnvConfig {
http_port,
webhook_secret,
open_router_api_key,
open_router_model,
open_router_timeout,
bot_max_concurrent,
gitea_url,
gitea_token,
gitea_timeout,
metrics_bind_addr,
})
}
pub fn try_get_env(key: &str) -> anyhow::Result<String> {
let env_value = std::env::var(key).map_err(|e| anyhow::anyhow!("{}: {}", key, e))?;
if env_value.trim().is_empty() {
return Err(anyhow!(format!("env var {} is empty", key)));
}
Ok(env_value)
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn test_try_get_env_returns_value() {
unsafe { std::env::set_var("TEST_ENV_PRESENT", "hello") };
assert_eq!(try_get_env("TEST_ENV_PRESENT").unwrap(), "hello");
}
#[test]
fn test_try_get_env_missing_var_returns_error() {
unsafe { std::env::remove_var("TEST_ENV_MISSING") };
assert!(try_get_env("TEST_ENV_MISSING").is_err());
}
#[test]
fn test_try_get_env_empty_value_returns_error() {
unsafe { std::env::set_var("TEST_ENV_EMPTY", "") };
let err = try_get_env("TEST_ENV_EMPTY").unwrap_err();
assert!(err.to_string().contains("TEST_ENV_EMPTY"));
}
#[test]
fn test_try_get_env_whitespace_only_returns_error() {
unsafe { std::env::set_var("TEST_ENV_WHITESPACE", " ") };
let err = try_get_env("TEST_ENV_WHITESPACE").unwrap_err();
assert!(err.to_string().contains("TEST_ENV_WHITESPACE"));
}
}
@@ -1,3 +1,4 @@
use anyhow::anyhow;
use axum::response::IntoResponse;
use reqwest::StatusCode;
@@ -21,6 +22,12 @@ pub enum AppError {
#[error("WebHook have bad action")]
InvalidActionErr,
#[error("Channel full")]
ChannelFullErr,
#[error("Already processed")]
AlreadyProcessedErr,
#[error(transparent)]
BadJsonStructErr(#[from] serde_json::Error),
@@ -54,10 +61,21 @@ impl IntoResponse for AppError {
StatusCode::UNAUTHORIZED,
"WebHook sig header is invalid".to_string(),
),
AppError::Other(_) => (
StatusCode::INTERNAL_SERVER_ERROR,
"Internal server error".to_string(),
),
AppError::AlreadyProcessedErr => (StatusCode::OK, "Already processed".to_string()),
AppError::ChannelFullErr => {
sentry_anyhow::capture_anyhow(&anyhow!("Max concurrent tasks reached"));
(
StatusCode::SERVICE_UNAVAILABLE,
"Max concurrent tasks reached".to_string(),
)
}
AppError::Other(err) => {
sentry_anyhow::capture_anyhow(&err);
(
StatusCode::INTERNAL_SERVER_ERROR,
"Internal server error".to_string(),
)
}
}
.into_response()
}
@@ -2,12 +2,11 @@ use std::time::Duration;
use serde::Deserialize;
use serde_json::{Value, json};
use tracing::instrument;
use crate::{
bot::{ReviewItem, ReviewResult},
errors::AppError,
};
use crate::{bot::ReviewResult, errors::AppError};
#[derive(Clone)]
pub struct GiteaAPI {
base_url: String,
client: reqwest::Client,
@@ -30,6 +29,22 @@ impl GiteaAPI {
})
}
#[instrument(skip(self))]
pub async fn get_authorized_user(&self) -> anyhow::Result<User> {
let url = format!("{}/api/v1/user", self.base_url);
let res = self.client.get(url).send().await?;
if !res.status().is_success() {
return Err(anyhow::anyhow!(
"Failed to get authorized user: {}",
res.status()
));
}
res.json::<User>().await.map_err(anyhow::Error::from)
}
#[instrument(skip(self))]
pub async fn comment(
&self,
body: &str,
@@ -57,6 +72,7 @@ impl GiteaAPI {
res.json::<Comment>().await.map_err(anyhow::Error::from)
}
#[instrument(skip(self))]
pub async fn edit_comment(
&self,
body: &str,
@@ -84,9 +100,30 @@ impl GiteaAPI {
Ok(())
}
#[instrument(skip(self))]
pub async fn delete_comment(&self, full_name: &str, comment_id: u64) -> anyhow::Result<()> {
let url = format!(
"{}/api/v1/repos/{}/issues/comments/{}",
self.base_url, full_name, comment_id
);
let res = self.client.delete(url).send().await?;
if !res.status().is_success() {
return Err(anyhow::anyhow!(
"Failed to delete comment: {}",
res.status()
));
}
Ok(())
}
#[instrument(skip(self, review_result))]
pub async fn post_pull_request_review(
&self,
review_result: &ReviewResult,
final_comment: &str,
full_name: &str,
index: u64,
) -> anyhow::Result<()> {
@@ -95,7 +132,7 @@ impl GiteaAPI {
self.base_url, full_name, index
);
let comments = &&review_result
let comments = &review_result
.reviews
.iter()
.filter(|r| r.line.is_some())
@@ -117,7 +154,7 @@ impl GiteaAPI {
.post(url)
.json(&json!({
"event": "COMMENT",
"body": review_result.comment,
"body": final_comment,
"comments": comments
}))
.send()
@@ -136,6 +173,20 @@ pub enum WebhookType {
Review(ReviewPayload),
}
impl WebhookType {
pub fn event_id(&self) -> u64 {
match self {
WebhookType::Review(payload) => payload.comment.id,
}
}
pub fn event_type_str(&self) -> &'static str {
match self {
WebhookType::Review(_) => "review",
}
}
}
#[derive(Deserialize, Debug)]
pub struct ReviewPayload {
pub action: String,
@@ -146,7 +197,6 @@ pub struct ReviewPayload {
#[derive(Deserialize, Debug)]
pub struct PullRequest {
pub id: u64,
pub diff_url: String,
pub number: u64,
pub title: String,
@@ -156,12 +206,11 @@ pub struct PullRequest {
pub struct Comment {
pub id: u64,
pub body: String,
pub user: User,
}
#[derive(Deserialize, Debug)]
pub struct User {
pub id: u64,
pub login: String,
}
#[derive(Deserialize, Debug)]
@@ -176,14 +225,6 @@ impl WebhookType {
_ => Err(AppError::UnknownEventErr),
}?;
let pr_body = match &wb {
WebhookType::Review(review_payload) => &review_payload.comment.body,
};
if !pr_body.starts_with(&format!("@{}", bot_name)) {
return Err(AppError::UnauthorizedUserErr);
}
let action = match &wb {
WebhookType::Review(review_payload) => &review_payload.action,
};
@@ -192,6 +233,14 @@ impl WebhookType {
return Err(AppError::InvalidActionErr);
}
let pr_body = match &wb {
WebhookType::Review(review_payload) => &review_payload.comment.body,
};
if !pr_body.starts_with(&format!("@{}", bot_name)) {
return Err(AppError::UnauthorizedUserErr);
}
Ok(wb)
}
}
@@ -207,13 +256,19 @@ mod tests {
"action": "created",
"pull_request": {
"id": 42,
"diff_url": "https://mydiff.fr"
"diff_url": "https://mydiff.fr",
"number": 1,
"title": "My PR"
},
"repository": {
"full_name": "owner/repo"
},
"comment": {
"id": 7,
"body": "@test_bot LGTM",
"user": {
"id": 100
"id": 100,
"login": "test_user"
}
}
});
@@ -224,10 +279,8 @@ mod tests {
match result.unwrap() {
WebhookType::Review(payload) => {
assert_eq!(payload.action, "created");
assert_eq!(payload.pull_request.id, 42);
assert_eq!(payload.comment.id, 7);
assert_eq!(payload.comment.body, "@test_bot LGTM");
assert_eq!(payload.comment.user.id, 100);
}
}
}
@@ -266,13 +319,19 @@ mod tests {
"action": "edited",
"pull_request": {
"id": 1,
"diff_url": "https://mydiff.fr"
"diff_url": "https://mydiff.fr",
"number": 1,
"title": "My PR"
},
"repository": {
"full_name": "owner/repo"
},
"comment": {
"id": 1,
"body": "@test_bot body",
"user": {
"id": 1
"id": 1,
"login": "test_user"
}
}
});
@@ -292,23 +351,27 @@ mod tests {
"action": "created",
"pull_request": {
"id": 99,
"diff_url": "https://mydiff.fr"
"diff_url": "https://mydiff.fr",
"number": 5,
"title": "My PR"
},
"repository": {
"full_name": "owner/repo"
},
"comment": {
"id": 12,
"body": "Needs work",
"user": {
"id": 200
"id": 200,
"login": "test_user"
}
}
});
let payload: ReviewPayload = serde_json::from_value(json).unwrap();
assert_eq!(payload.action, "created");
assert_eq!(payload.pull_request.id, 99);
assert_eq!(payload.comment.id, 12);
assert_eq!(payload.comment.body, "Needs work");
assert_eq!(payload.comment.user.id, 200);
}
#[test]
@@ -324,13 +387,19 @@ mod tests {
"action": "created",
"pull_request": {
"id": 1,
"diff_url": "https://mydiff.fr"
"diff_url": "https://mydiff.fr",
"number": 1,
"title": "My PR"
},
"repository": {
"full_name": "owner/repo"
},
"comment": {
"id": 1,
"body": "@other_bot do something",
"user": {
"id": 1
"id": 1,
"login": "test_user"
}
}
});
@@ -345,13 +414,19 @@ mod tests {
"action": "created",
"pull_request": {
"id": 1,
"diff_url": "https://mydiff.fr"
"diff_url": "https://mydiff.fr",
"number": 1,
"title": "My PR"
},
"repository": {
"full_name": "owner/repo"
},
"comment": {
"id": 1,
"body": "just a comment without bot mention",
"user": {
"id": 1
"id": 1,
"login": "test_user"
}
}
});
+119
View File
@@ -0,0 +1,119 @@
use crate::{
bot::Bot,
gitea::{GiteaAPI, WebhookType},
open_router::OpenRouterClient,
state::AppState,
};
use dotenvy::dotenv;
use tokio::signal::unix::{SignalKind, signal};
use tokio_util::sync::CancellationToken;
use tracing::{info, warn};
use tracing_subscriber::{EnvFilter, fmt, layer::SubscriberExt, util::SubscriberInitExt};
mod api;
mod bot;
mod bot_actions;
mod consts;
mod env;
mod errors;
mod gitea;
mod metrics;
mod open_router;
mod state;
fn main() -> anyhow::Result<()> {
dotenv().ok();
tracing_subscriber::registry()
.with(fmt::layer())
.with(
EnvFilter::try_from_default_env() // lit RUST_LOG depuis l'env
.unwrap_or_else(|_| EnvFilter::new("info")),
)
.init();
let _sentry_guard = if let Ok(sentry_dsn) = env::try_get_env("SENTRY_DSN") {
info!("Initialize sentry");
Some(sentry::init((
sentry_dsn,
sentry::ClientOptions {
release: sentry::release_name!(),
send_default_pii: true,
..Default::default()
},
)))
} else {
warn!("SENTRY_DSN not set, sentry will not be initialized");
None
};
tokio::runtime::Runtime::new()?.block_on(run())
}
async fn run() -> anyhow::Result<()> {
let config = env::load_config()?;
if let Some(metric_bind_addr) = &config.metrics_bind_addr {
metrics::install(metric_bind_addr)?;
}
let gitea_api = GiteaAPI::new(&config.gitea_url, &config.gitea_token, config.gitea_timeout)?;
let gitea_user = gitea_api.get_authorized_user().await?;
info!(
port = config.http_port,
model = %config.open_router_model,
gitea_url = %config.gitea_url,
bot_name = %gitea_user.login,
"Starting Herald"
);
let open_router_client = OpenRouterClient::new(
&config.open_router_api_key,
&config.open_router_model,
config.open_router_timeout,
)?;
let shutdown = CancellationToken::new();
let bot = Bot::new(
gitea_user.login,
gitea_api,
open_router_client,
reqwest::Client::new(),
config.bot_max_concurrent,
config.open_router_model.clone(),
);
let (tx, rx) = tokio::sync::mpsc::channel::<WebhookType>(config.bot_max_concurrent * 2);
let app_state = AppState {
bot_tx: tx,
bot: bot.clone(),
config,
};
let signal = async {
let mut sigterm = signal(SignalKind::terminate())?;
let mut sigint = signal(SignalKind::interrupt())?;
tokio::select! {
_ = sigterm.recv() => info!("Received SIGTERM"),
_ = sigint.recv() => info!("Received SIGINT"),
}
info!("Shutting down...");
shutdown.cancel();
anyhow::Ok(())
};
tokio::try_join!(
bot.start(rx, shutdown.clone()),
api::start(app_state, shutdown.clone()),
signal
)?;
info!("Shutdown complete");
Ok(())
}
+88
View File
@@ -0,0 +1,88 @@
use std::{net::SocketAddr, str::FromStr};
use metrics::{Unit, counter, describe_counter, describe_gauge, gauge};
pub fn webhook_received(event_type: &str) {
counter!("herald_webhooks_received_total", "event_type" => event_type.to_string()).increment(1);
}
pub fn webhook_duplicate(event_type: &str) {
counter!("herald_webhooks_duplicate_total", "event_type" => event_type.to_string())
.increment(1);
}
pub fn webhook_channel_full(event_type: &str) {
counter!("herald_webhooks_channel_full_total", "event_type" => event_type.to_string())
.increment(1);
}
pub fn increment_task_active() {
gauge!("herald_bot_tasks_active").increment(1.0);
}
pub fn decrement_task_active() {
gauge!("herald_bot_tasks_active").decrement(1.0);
}
pub fn task_completed(event_type: &str) {
counter!("herald_bot_tasks_completed_total", "event_type" => event_type.to_string())
.increment(1);
}
pub fn task_failed(event_type: &str) {
counter!("herald_bot_tasks_failed_total", "event_type" => event_type.to_string()).increment(1);
}
pub fn openrouter_cost_usd(cost: f64) {
counter!("herald_openrouter_cost_cents_total").increment((cost * 100.0).round() as u64);
}
pub fn describe() {
describe_counter!(
"herald_webhooks_received_total",
Unit::Count,
"Total webhooks received"
);
describe_counter!(
"herald_webhooks_duplicate_total",
Unit::Count,
"Webhooks rejected as duplicates"
);
describe_counter!(
"herald_webhooks_channel_full_total",
Unit::Count,
"Webhooks dropped because the bot channel was full"
);
describe_gauge!(
"herald_bot_tasks_active",
Unit::Count,
"Bot tasks currently in progress"
);
describe_counter!(
"herald_bot_tasks_completed_total",
Unit::Count,
"Bot tasks completed successfully"
);
describe_counter!(
"herald_bot_tasks_failed_total",
Unit::Count,
"Bot tasks that failed"
);
describe_counter!(
"herald_openrouter_cost_cents_total",
Unit::Count,
"Total OpenRouter cost in cents (divide by 100 for USD)"
);
}
pub fn install(bind_addr: &str) -> anyhow::Result<()> {
describe();
let builder = metrics_exporter_prometheus::PrometheusBuilder::new();
builder
.with_http_listener(SocketAddr::from_str(bind_addr)?)
.install()?;
tracing::info!(bind_addr, "Prometheus metrics exporter installed");
Ok(())
}
@@ -1,12 +1,14 @@
use std::time::Duration;
use openrouter_rs::{Message, api::chat::ChatCompletionRequest};
use tracing::instrument;
pub struct ChatResult {
pub message: String,
pub cost: Option<f64>,
}
#[derive(Clone)]
pub struct OpenRouterClient {
client: openrouter_rs::OpenRouterClient,
model: String,
@@ -27,6 +29,7 @@ impl OpenRouterClient {
})
}
#[instrument(skip(self), err)]
pub async fn chat(&self, msg: &str) -> anyhow::Result<ChatResult> {
let request = ChatCompletionRequest::builder()
.model(&self.model)
@@ -39,7 +42,7 @@ impl OpenRouterClient {
Ok(ChatResult {
message: response.choices[0]
.content()
.map(|msg| String::from(msg))
.map(String::from)
.ok_or(anyhow::anyhow!("No content"))?,
cost: response.usage.and_then(|u| u.cost),
})
@@ -1,7 +1,8 @@
use crate::{env::EnvConfig, gitea::WebhookType};
use crate::{bot::Bot, env::EnvConfig, gitea::WebhookType};
#[derive(Clone)]
pub struct AppState {
pub bot_tx: tokio::sync::mpsc::Sender<WebhookType>,
pub bot: Bot,
pub config: EnvConfig,
}
+22
View File
@@ -0,0 +1,22 @@
#!/usr/bin/env bash
set -euo pipefail
IMAGE="tintounn/herald"
if [ -z "${CI_COMMIT_TAG:-}" ]; then
echo "Error: CI_COMMIT_TAG is not set" >&2
exit 1
fi
TAG="${CI_COMMIT_TAG}"
echo "Building ${IMAGE}:${TAG}..."
buildah build \
--file Containerfile \
--tag "docker.io/${IMAGE}:${TAG}" \
.
echo "Pushing ${IMAGE}:${TAG}..."
buildah push "${IMAGE}:${TAG}"
echo "Done: ${IMAGE}:${TAG}"
-117
View File
@@ -1,117 +0,0 @@
use serde::Deserialize;
use std::time::Duration;
use crate::{
env::EnvConfig,
gitea::{GiteaAPI, WebhookType},
open_router::OpenRouterClient,
};
#[derive(Deserialize, Debug)]
pub struct ReviewResult {
pub reviews: Vec<ReviewItem>,
pub comment: String,
pub cost: Option<f64>,
}
#[derive(Deserialize, Debug)]
pub struct ReviewItem {
pub filename: String,
pub line: Option<u64>,
pub code: String,
pub message: String,
}
/// Map a filename to a markdown language identifier for syntax highlighting.
fn lang_from_filename(filename: &str) -> &str {
match std::path::Path::new(filename)
.extension()
.and_then(|e| e.to_str())
.unwrap_or("")
{
"rs" => "rust",
"py" => "python",
"js" | "mjs" => "javascript",
"ts" => "typescript",
"jsx" => "jsx",
"tsx" => "tsx",
"go" => "go",
"java" => "java",
"kt" | "kts" => "kotlin",
"scala" => "scala",
"c" | "h" => "c",
"cpp" | "cc" | "cxx" | "hpp" | "hxx" => "cpp",
"rb" => "ruby",
"php" => "php",
"swift" => "swift",
"sh" | "bash" | "zsh" => "bash",
"sql" => "sql",
"html" | "htm" => "html",
"css" => "css",
"scss" | "sass" => "scss",
"json" => "json",
"yaml" | "yml" => "yaml",
"xml" => "xml",
"toml" => "toml",
"md" | "mdx" => "markdown",
"dockerfile" | "Dockerfile" => "dockerfile",
"Makefile" => "makefile",
_ => "",
}
}
pub struct Bot {
config: EnvConfig,
gitea_api: GiteaAPI,
open_router_client: OpenRouterClient,
http_client: reqwest::Client,
}
impl Bot {
pub fn new(config: EnvConfig) -> anyhow::Result<Self> {
let gitea_timeout = config.gitea_timeout;
let open_router_timeout = config.open_router_timeout;
Ok(Self {
gitea_api: GiteaAPI::new(&config.gitea_url, &config.gitea_token, gitea_timeout)?,
open_router_client: OpenRouterClient::new(
&config.open_router_api_key,
&config.open_router_model,
open_router_timeout,
)?,
config,
http_client: reqwest::Client::builder()
.timeout(Duration::from_secs(gitea_timeout))
.build()?,
})
}
pub async fn start(
&self,
mut rx: tokio::sync::mpsc::Receiver<WebhookType>,
) -> anyhow::Result<()> {
while let Some(wb) = rx.recv().await {
self.exec(wb).await;
}
Ok(())
}
pub async fn exec(&self, webhook: WebhookType) {
let exec_result = match webhook {
WebhookType::Review(review_payload) => crate::bot_actions::review::exec_review(
&self.gitea_api,
&self.open_router_client,
&self.http_client,
&self.config.open_router_model,
review_payload,
),
}
.await;
match exec_result {
Ok(_) => println!("Task completed"),
Err(err) => println!("{}", err),
}
}
}
-51
View File
@@ -1,51 +0,0 @@
use anyhow::anyhow;
use dotenvy::dotenv;
#[derive(Clone)]
pub struct EnvConfig {
pub http_port: u16,
pub webhook_secret: String,
pub open_router_api_key: String,
pub open_router_model: String,
pub open_router_timeout: u64,
pub bot_name: String,
pub gitea_url: String,
pub gitea_token: String,
pub gitea_timeout: u64,
}
pub fn load_config() -> anyhow::Result<EnvConfig> {
dotenv().ok();
let http_port = try_get_env("HTTP_PORT")?.parse()?;
let bot_name = try_get_env("BOT_NAME")?;
let webhook_secret = try_get_env("WEBHOOK_SIG_HEADER_SECRET")?;
let open_router_api_key = try_get_env("OPEN_ROUTER_API_KEY")?;
let open_router_model = try_get_env("OPEN_ROUTER_MODEL")?;
let open_router_timeout = try_get_env("OPEN_ROUTER_TIMEOUT")?.parse()?;
let gitea_url = try_get_env("GITEA_URL")?;
let gitea_token = try_get_env("GITEA_TOKEN")?;
let gitea_timeout = try_get_env("GITEA_TIMEOUT")?.parse()?;
Ok(EnvConfig {
http_port,
webhook_secret,
bot_name,
open_router_api_key,
open_router_model,
open_router_timeout,
gitea_url,
gitea_token,
gitea_timeout,
})
}
fn try_get_env(key: &str) -> anyhow::Result<String> {
let env = std::env::var(key)?;
if env.trim().is_empty() {
return Err(anyhow!(format!("env var {} is empty", key)));
}
Ok(env)
}
-25
View File
@@ -1,25 +0,0 @@
use crate::{bot::Bot, gitea::WebhookType, state::AppState};
mod api;
mod bot;
mod bot_actions;
mod consts;
mod env;
mod errors;
mod gitea;
mod open_router;
mod state;
#[tokio::main]
async fn main() -> anyhow::Result<()> {
let config = env::load_config()?;
let bot = Bot::new(config.clone())?;
let (tx, rx) = tokio::sync::mpsc::channel::<WebhookType>(1);
let app_state = AppState { bot_tx: tx, config };
tokio::try_join!(bot.start(rx), api::start(app_state))?;
Ok(())
}