fixed typos for redis
This commit is contained in:
@@ -303,7 +303,7 @@ impl Database {
|
|||||||
error
|
error
|
||||||
)
|
)
|
||||||
VALUES (
|
VALUES (
|
||||||
TO_TIMESTAMP($1::DOUBLE PRECISION / 1000.0),
|
TO_TIMESTAMP(($1::BIGINT)::DOUBLE PRECISION / 1000.0),
|
||||||
$1,
|
$1,
|
||||||
$2,
|
$2,
|
||||||
$3,
|
$3,
|
||||||
|
|||||||
+49
-3
@@ -43,6 +43,23 @@ pub async fn run_redis_coordinator(
|
|||||||
let result_queue = result_queue_name(&config.redis_queue_name);
|
let result_queue = result_queue_name(&config.redis_queue_name);
|
||||||
let total_key = total_ips_key(&config.redis_queue_name);
|
let total_key = total_ips_key(&config.redis_queue_name);
|
||||||
let completed_key = completed_ips_key(&config.redis_queue_name);
|
let completed_key = completed_ips_key(&config.redis_queue_name);
|
||||||
|
let should_resume = should_resume_existing_run(
|
||||||
|
&config.redis_url,
|
||||||
|
&config.redis_queue_name,
|
||||||
|
&result_queue,
|
||||||
|
&total_key,
|
||||||
|
&completed_key,
|
||||||
|
ips.len(),
|
||||||
|
)
|
||||||
|
.await?;
|
||||||
|
|
||||||
|
if should_resume {
|
||||||
|
println!(
|
||||||
|
"Retomando fila Redis existente em {}. Total esperado: {} IPs",
|
||||||
|
config.redis_queue_name,
|
||||||
|
ips.len()
|
||||||
|
);
|
||||||
|
} else {
|
||||||
clear_queue(
|
clear_queue(
|
||||||
&config.redis_url,
|
&config.redis_url,
|
||||||
&config.redis_queue_name,
|
&config.redis_queue_name,
|
||||||
@@ -80,6 +97,7 @@ pub async fn run_redis_coordinator(
|
|||||||
config.redis_worker_count.max(1)
|
config.redis_worker_count.max(1)
|
||||||
);
|
);
|
||||||
}
|
}
|
||||||
|
}
|
||||||
|
|
||||||
consume_result_queue(
|
consume_result_queue(
|
||||||
&config.redis_url,
|
&config.redis_url,
|
||||||
@@ -708,6 +726,28 @@ async fn read_global_progress(
|
|||||||
})
|
})
|
||||||
}
|
}
|
||||||
|
|
||||||
|
async fn should_resume_existing_run(
|
||||||
|
redis_url: &str,
|
||||||
|
queue_name: &str,
|
||||||
|
result_queue: &str,
|
||||||
|
total_key: &str,
|
||||||
|
completed_key: &str,
|
||||||
|
expected_total: usize,
|
||||||
|
) -> Result<bool> {
|
||||||
|
let existing_total = redis_usize(redis_url, total_key).await?.unwrap_or_default();
|
||||||
|
if existing_total != expected_total {
|
||||||
|
return Ok(false);
|
||||||
|
}
|
||||||
|
|
||||||
|
let completed = redis_usize(redis_url, completed_key)
|
||||||
|
.await?
|
||||||
|
.unwrap_or_default();
|
||||||
|
let queued_work = redis_list_len(redis_url, queue_name).await?;
|
||||||
|
let queued_results = redis_list_len(redis_url, result_queue).await?;
|
||||||
|
|
||||||
|
Ok(queued_work > 0 || queued_results > 0 || (completed > 0 && completed < expected_total))
|
||||||
|
}
|
||||||
|
|
||||||
fn total_ips_key(queue_name: &str) -> String {
|
fn total_ips_key(queue_name: &str) -> String {
|
||||||
format!("{queue_name}:total_ips")
|
format!("{queue_name}:total_ips")
|
||||||
}
|
}
|
||||||
@@ -730,8 +770,8 @@ async fn push_fetch_result(
|
|||||||
Ok(())
|
Ok(())
|
||||||
}
|
}
|
||||||
|
|
||||||
async fn pop_fetch_result(redis_url: &str, result_queue: &str) -> Result<Option<RtspFetchResult>> {
|
async fn peek_fetch_result(redis_url: &str, result_queue: &str) -> Result<Option<RtspFetchResult>> {
|
||||||
let output = redis_cli(redis_url, &["--raw", "LPOP", result_queue]).await?;
|
let output = redis_cli(redis_url, &["--raw", "LINDEX", result_queue, "0"]).await?;
|
||||||
let value = output.trim_end_matches('\n');
|
let value = output.trim_end_matches('\n');
|
||||||
if value.is_empty() {
|
if value.is_empty() {
|
||||||
return Ok(None);
|
return Ok(None);
|
||||||
@@ -739,6 +779,11 @@ async fn pop_fetch_result(redis_url: &str, result_queue: &str) -> Result<Option<
|
|||||||
deserialize_fetch_result(value).map(Some)
|
deserialize_fetch_result(value).map(Some)
|
||||||
}
|
}
|
||||||
|
|
||||||
|
async fn ack_fetch_result(redis_url: &str, result_queue: &str) -> Result<()> {
|
||||||
|
let _ = redis_cli(redis_url, &["LPOP", result_queue]).await?;
|
||||||
|
Ok(())
|
||||||
|
}
|
||||||
|
|
||||||
async fn consume_result_queue(
|
async fn consume_result_queue(
|
||||||
redis_url: &str,
|
redis_url: &str,
|
||||||
queue_name: &str,
|
queue_name: &str,
|
||||||
@@ -750,7 +795,7 @@ async fn consume_result_queue(
|
|||||||
let mut inserted = 0usize;
|
let mut inserted = 0usize;
|
||||||
loop {
|
loop {
|
||||||
let mut drained_any = false;
|
let mut drained_any = false;
|
||||||
while let Some(result) = pop_fetch_result(redis_url, result_queue).await? {
|
while let Some(result) = peek_fetch_result(redis_url, result_queue).await? {
|
||||||
database
|
database
|
||||||
.insert_rtsp_fetch_result(
|
.insert_rtsp_fetch_result(
|
||||||
result.checked_at,
|
result.checked_at,
|
||||||
@@ -763,6 +808,7 @@ async fn consume_result_queue(
|
|||||||
result.error.as_deref(),
|
result.error.as_deref(),
|
||||||
)
|
)
|
||||||
.await?;
|
.await?;
|
||||||
|
ack_fetch_result(redis_url, result_queue).await?;
|
||||||
inserted += 1;
|
inserted += 1;
|
||||||
drained_any = true;
|
drained_any = true;
|
||||||
}
|
}
|
||||||
|
|||||||
Reference in New Issue
Block a user