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

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
8 changes: 4 additions & 4 deletions Cargo.lock

Some generated files are not rendered by default. Learn more about how customized files appear on GitHub.

2 changes: 1 addition & 1 deletion Cargo.toml
Original file line number Diff line number Diff line change
Expand Up @@ -10,7 +10,7 @@ members = [
resolver = "2"

[workspace.package]
version = "0.1.7"
version = "0.1.8"

[workspace.dependencies]
anyhow = "1.0.102"
Expand Down
12 changes: 8 additions & 4 deletions cuscuta-entry/src/main.rs
Original file line number Diff line number Diff line change
Expand Up @@ -40,7 +40,8 @@ use reqwest::StatusCode;
use serde_json::json;
use tokio::net::TcpListener;
use tokio_util::sync::CancellationToken;
use tower_http::trace::TraceLayer;
use tower_http::trace::{self, TraceLayer};
use tracing::Level;

use crate::{
endpoints::enqueue,
Expand All @@ -55,12 +56,15 @@ async fn main() {
.init();
tracing::info!("starting...");
let halt_token = CancellationToken::new();
let trace_layer = TraceLayer::new_for_http()
.make_span_with(trace::DefaultMakeSpan::new().level(Level::INFO))
.on_request(trace::DefaultOnRequest::new().level(Level::INFO))
.on_response(trace::DefaultOnResponse::new().level(Level::INFO));
let service = Router::new()
.route("/healthz", get(healthz))
.route("/readyz", get(readyz))
.route("/v1/enqueue", post(enqueue))
.route("/v1/query", get(query))
.layer(TraceLayer::new_for_http());
.route("/v1/enqueue", post(enqueue).layer(trace_layer.clone()))
.route("/v1/query", get(query).layer(trace_layer));
let addr = TcpListener::bind("0.0.0.0:8081")
.await
.expect("failed to bind 0.0.0.0:8081");
Expand Down
9 changes: 7 additions & 2 deletions cuscuta-worker/src/worker/clean.rs
Original file line number Diff line number Diff line change
Expand Up @@ -90,7 +90,7 @@ pub async fn clean_jobs(
&& !pending_friends_code.contains(&finished_job.essential.friend_code)
{
let friend_user_id = friend_info.user_id.to_string();
xxxxxx_safe_call(
if let Err(e) = xxxxxx_safe_call(
config.worker_max_retry_count,
config.worker_exponential_backoff_base_millis,
config.worker_exponential_backoff_multiplier,
Expand All @@ -106,7 +106,12 @@ pub async fn clean_jobs(
},
)
.await
.map_err(Error::Api)?;
{
worker_write_event!(
WorkerEventType::Warn,
format!("failed to delete friend: {e}")
);
}
friends.retain(|it| it.user_id != friend_info.user_id);
}
let cursor_length = i64::from(finished_job.essential.cursor_length);
Expand Down
2 changes: 1 addition & 1 deletion cuscuta-worker/src/worker/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -79,6 +79,7 @@ pub async fn worker_loop(cancellation_token: &CancellationToken) -> WorkerResult
WORKER_ID.get_or_init(|| worker_id.clone());
while !cancellation_token.is_cancelled() {
if let Err(e) = internal_loop(&mut current_jobs, &mut cursor, &mut friends).await {
worker_write_event!(WorkerEventType::Warn, format!("worker loop failed: {e}"));
if let Error::Api(api_error) = &e {
match api_error {
api::Error::Network(_) => {}
Expand All @@ -92,7 +93,6 @@ pub async fn worker_loop(cancellation_token: &CancellationToken) -> WorkerResult
}
}
}
worker_write_event!(WorkerEventType::Warn, format!("worker loop failed: {e}"));
sleep(Duration::from_secs(1)).await;
}
}
Expand Down
Loading