mirror of
https://github.com/Portabase/agent.git
synced 2026-09-11 02:27:10 +00:00
Compare commits
4 Commits
| Author | SHA1 | Date | |
|---|---|---|---|
| 4df3fe2d0c | |||
| ec7715ca26 | |||
| 1d15f40662 | |||
| 324a2ea3f2 |
+1
-1
@@ -27,5 +27,5 @@ keywords:
|
||||
- self-hosted
|
||||
- portabase
|
||||
license: Apache-2.0
|
||||
version: 1.16.0
|
||||
version: 1.16.2
|
||||
date-released: '2026-02-24'
|
||||
|
||||
Generated
+1
-1
@@ -3503,7 +3503,7 @@ dependencies = [
|
||||
|
||||
[[package]]
|
||||
name = "portabase-agent"
|
||||
version = "1.16.0"
|
||||
version = "1.16.2"
|
||||
dependencies = [
|
||||
"aes",
|
||||
"aes-gcm",
|
||||
|
||||
+1
-1
@@ -1,6 +1,6 @@
|
||||
[package]
|
||||
name = "portabase-agent"
|
||||
version = "1.16.0"
|
||||
version = "1.16.2"
|
||||
edition = "2024"
|
||||
|
||||
[dependencies]
|
||||
|
||||
+1
-1
@@ -19,7 +19,7 @@ services:
|
||||
APP_ENV: development
|
||||
LOG: debug
|
||||
TZ: "Europe/Paris"
|
||||
EDGE_KEY: "eyJzZXJ2ZXJVcmwiOiJodHRwOi8vbG9jYWxob3N0Ojg4ODciLCJhZ2VudElkIjoiNDA1MTA4YzQtMDJjYy00NTlhLTkxNjItODExNTc3NjAzMjhjIiwibWFzdGVyS2V5QjY0IjoiMUh0djdtWCtYVkJxL0IzUEV2WDlZZjlQeUdVZW5oRHlXemo5THRqNW90WT0ifQ=="
|
||||
EDGE_KEY: "eyJzZXJ2ZXJVcmwiOiJodHRwOi8vbG9jYWxob3N0Ojg4ODciLCJhZ2VudElkIjoiMWM4NmQ5NGEtMGVjOC00NzkxLTk0ZTEtOWRlYmJmMWY0M2I5IiwibWFzdGVyS2V5QjY0IjoiMUh0djdtWCtYVkJxL0IzUEV2WDlZZjlQeUdVZW5oRHlXemo5THRqNW90WT0ifQ=="
|
||||
#CHUNK_SIZE_MB: "1"
|
||||
#POOLING: 1
|
||||
#DATABASES_CONFIG_FILE: "config.toml"
|
||||
|
||||
@@ -104,6 +104,7 @@ impl BackupService {
|
||||
method,
|
||||
&storage,
|
||||
Some(encrypt),
|
||||
&backup_storage_id,
|
||||
)
|
||||
.await;
|
||||
|
||||
|
||||
@@ -22,6 +22,7 @@ pub trait StorageProvider: Send + Sync {
|
||||
method: BackupMethod,
|
||||
config: &DatabaseStorage,
|
||||
encrypt: Option<bool>,
|
||||
backup_storage_id: &str,
|
||||
) -> UploadResult;
|
||||
}
|
||||
|
||||
|
||||
@@ -26,6 +26,7 @@ impl StorageProvider for AzureBlobProvider {
|
||||
_method: BackupMethod,
|
||||
storage: &DatabaseStorage,
|
||||
encrypt: Option<bool>,
|
||||
_backup_storage_id: &str,
|
||||
) -> UploadResult {
|
||||
let Some(file_path) = result.backup_file else {
|
||||
return UploadResult {
|
||||
|
||||
@@ -28,6 +28,7 @@ impl StorageProvider for GoogleCloudStorageProvider {
|
||||
_method: BackupMethod,
|
||||
storage: &DatabaseStorage,
|
||||
encrypt: Option<bool>,
|
||||
_backup_storage_id: &str,
|
||||
) -> UploadResult {
|
||||
let Some(file_path) = result.backup_file else {
|
||||
return UploadResult {
|
||||
|
||||
@@ -26,6 +26,7 @@ impl StorageProvider for GoogleDriveProvider {
|
||||
_method: BackupMethod,
|
||||
storage: &DatabaseStorage,
|
||||
encrypt: Option<bool>,
|
||||
_backup_storage_id: &str,
|
||||
) -> UploadResult {
|
||||
let Some(file_path) = result.backup_file else {
|
||||
return UploadResult {
|
||||
|
||||
@@ -23,6 +23,7 @@ impl StorageProvider for LocalProvider {
|
||||
method: BackupMethod,
|
||||
storage: &DatabaseStorage,
|
||||
encrypt: Option<bool>,
|
||||
backup_storage_id: &str,
|
||||
) -> UploadResult {
|
||||
let Some(file_path) = result.backup_file else {
|
||||
return UploadResult {
|
||||
@@ -82,6 +83,10 @@ impl StorageProvider for LocalProvider {
|
||||
"X-Generated-Id",
|
||||
HeaderValue::from_str(&result.generated_id).unwrap(),
|
||||
);
|
||||
extra_headers.insert(
|
||||
"X-Backup-Storage-Id",
|
||||
HeaderValue::from_str(backup_storage_id).unwrap(),
|
||||
);
|
||||
extra_headers.insert("X-Status", HeaderValue::from_str(&result.status).unwrap());
|
||||
extra_headers.insert(
|
||||
"X-Method",
|
||||
|
||||
@@ -34,6 +34,7 @@ impl StorageProvider for S3Provider {
|
||||
_method: BackupMethod,
|
||||
storage: &DatabaseStorage,
|
||||
encrypt: Option<bool>,
|
||||
_backup_storage_id: &str,
|
||||
) -> UploadResult {
|
||||
let Some(file_path) = result.backup_file else {
|
||||
return UploadResult {
|
||||
|
||||
@@ -31,12 +31,45 @@ fn normalized_expression_is_valid_for_cron_schedule() {
|
||||
#[test]
|
||||
fn next_run_timestamp_returns_future_timestamp() {
|
||||
let expr = normalize_cron("*/1 * * * *");
|
||||
let ts = next_run_timestamp(&expr);
|
||||
let ts = next_run_timestamp(&expr).unwrap();
|
||||
|
||||
let now = chrono::Local::now().timestamp();
|
||||
assert!(ts > now);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn normalize_converts_unix_sunday_zero_to_crate_dow() {
|
||||
let input = "00 06 * * 0";
|
||||
let normalized = normalize_cron(input);
|
||||
|
||||
assert_eq!(normalized, "0 00 06 * * 1");
|
||||
|
||||
let schedule = Schedule::from_str(&normalized);
|
||||
assert!(schedule.is_ok());
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn normalize_converts_unix_saturday_to_crate_dow() {
|
||||
assert_eq!(normalize_cron("00 06 * * 6"), "0 00 06 * * 7");
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn normalize_converts_unix_dow_range() {
|
||||
assert_eq!(normalize_cron("00 06 * * 0-4"), "0 00 06 * * 1-5");
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn next_run_timestamp_returns_none_for_invalid_cron() {
|
||||
assert!(next_run_timestamp("not a cron").is_none());
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn normalize_leaves_out_of_range_dow_untouched() {
|
||||
let normalized = normalize_cron("* * * * 100");
|
||||
assert_eq!(normalized, "0 * * * * 100");
|
||||
assert!(Schedule::from_str(&normalized).is_err());
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn normalization_does_not_break_schedule_parsing() {
|
||||
let input = "0 */10 * * * *";
|
||||
|
||||
@@ -10,9 +10,9 @@ use std::str::FromStr;
|
||||
use tracing::debug;
|
||||
use tracing::info;
|
||||
|
||||
pub fn next_run_timestamp(expr: &str) -> i64 {
|
||||
let schedule = Schedule::from_str(expr).unwrap();
|
||||
schedule.upcoming(Local).next().unwrap().timestamp()
|
||||
pub fn next_run_timestamp(expr: &str) -> Option<i64> {
|
||||
let schedule = Schedule::from_str(expr).ok()?;
|
||||
Some(schedule.upcoming(Local).next()?.timestamp())
|
||||
}
|
||||
|
||||
pub async fn check_and_update_cron(
|
||||
@@ -50,24 +50,21 @@ pub async fn check_and_update_cron(
|
||||
let metadata_changed = stored.metadata != metadata;
|
||||
|
||||
if cron_changed || args_changed || metadata_changed {
|
||||
upsert_task(conn, &task_name, task, &cron, args.clone(), metadata)
|
||||
.await
|
||||
.unwrap_or_else(|e| {
|
||||
tracing::error!("Failed to update task {}: {:?}", task_name, e);
|
||||
});
|
||||
|
||||
info!(
|
||||
"Task {} updated (cron: {}, args: {}, metadata: {})",
|
||||
task_name, cron_changed, args_changed, metadata_changed
|
||||
);
|
||||
match upsert_task(conn, &task_name, task, &cron, args.clone(), metadata).await {
|
||||
Ok(()) => info!(
|
||||
"Task {} updated (cron: {}, args: {}, metadata: {})",
|
||||
task_name, cron_changed, args_changed, metadata_changed
|
||||
),
|
||||
Err(e) => {
|
||||
tracing::error!("Failed to update task {}: {:?}", task_name, e)
|
||||
}
|
||||
}
|
||||
}
|
||||
} else {
|
||||
upsert_task(conn, &task_name, task, &cron, args, metadata)
|
||||
.await
|
||||
.unwrap_or_else(|e| {
|
||||
tracing::error!("Failed to create task {}: {:?}", task_name, e);
|
||||
});
|
||||
info!("Task {} created", task_name);
|
||||
match upsert_task(conn, &task_name, task, &cron, args, metadata).await {
|
||||
Ok(()) => info!("Task {} created", task_name),
|
||||
Err(e) => tracing::error!("Failed to create task {}: {:?}", task_name, e),
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
@@ -48,8 +48,17 @@ pub async fn scheduler_loop(mut conn: MultiplexedConnection) {
|
||||
task_clone.task, e
|
||||
);
|
||||
}
|
||||
let next_ts = next_run_timestamp(&task_clone.cron);
|
||||
let _: () = conn_clone.zadd(SCHEDULE_KEY, &key, next_ts).await.unwrap();
|
||||
match next_run_timestamp(&task_clone.cron) {
|
||||
Some(next_ts) => {
|
||||
let _: () = conn_clone.zadd(SCHEDULE_KEY, &key, next_ts).await.unwrap();
|
||||
}
|
||||
None => {
|
||||
error!(
|
||||
"Invalid cron expression for task={}: {}",
|
||||
task_clone.task, task_clone.cron
|
||||
);
|
||||
}
|
||||
}
|
||||
});
|
||||
}
|
||||
tokio::time::sleep(std::time::Duration::from_secs(1)).await;
|
||||
|
||||
@@ -16,7 +16,16 @@ pub async fn upsert_task(
|
||||
metadata: Option<Value>,
|
||||
) -> redis::RedisResult<()> {
|
||||
let key = format!("redbeat:{}", name);
|
||||
let next_ts = next_run_timestamp(cron);
|
||||
let next_ts = match next_run_timestamp(cron) {
|
||||
Some(ts) => ts,
|
||||
None => {
|
||||
return Err(redis::RedisError::from((
|
||||
redis::ErrorKind::Client,
|
||||
"invalid cron expression",
|
||||
cron.to_string(),
|
||||
)));
|
||||
}
|
||||
};
|
||||
|
||||
let entry = PeriodicTask {
|
||||
task: task.to_string(),
|
||||
|
||||
+51
-3
@@ -1,7 +1,55 @@
|
||||
pub fn normalize_cron(expr: &str) -> String {
|
||||
if expr.split_whitespace().count() == 5 {
|
||||
format!("0 {}", expr)
|
||||
let fields: Vec<&str> = expr.split_whitespace().collect();
|
||||
|
||||
let (sec, mut rest): (String, Vec<String>) = match fields.len() {
|
||||
5 => ("0".to_string(), fields.iter().map(|s| s.to_string()).collect()),
|
||||
6 => (
|
||||
fields[0].to_string(),
|
||||
fields[1..].iter().map(|s| s.to_string()).collect(),
|
||||
),
|
||||
_ => return expr.to_string(),
|
||||
};
|
||||
|
||||
if let Some(last) = rest.last_mut() {
|
||||
*last = convert_dow(last);
|
||||
}
|
||||
|
||||
format!("{} {}", sec, rest.join(" "))
|
||||
}
|
||||
|
||||
fn convert_dow(field: &str) -> String {
|
||||
field
|
||||
.split(',')
|
||||
.map(convert_dow_part)
|
||||
.collect::<Vec<_>>()
|
||||
.join(",")
|
||||
}
|
||||
|
||||
fn convert_dow_part(part: &str) -> String {
|
||||
let (base, step) = match part.split_once('/') {
|
||||
Some((b, s)) => (b, Some(s)),
|
||||
None => (part, None),
|
||||
};
|
||||
|
||||
let converted = if let Some((start, end)) = base.split_once('-') {
|
||||
match (start.parse::<u8>(), end.parse::<u8>()) {
|
||||
(Ok(a), Ok(b)) if a <= 7 && b <= 7 => {
|
||||
format!("{}-{}", (a % 7) + 1, (b % 7) + 1)
|
||||
}
|
||||
_ => base.to_string(),
|
||||
}
|
||||
} else if let Ok(n) = base.parse::<u8>() {
|
||||
if n <= 7 {
|
||||
((n % 7) + 1).to_string()
|
||||
} else {
|
||||
base.to_string()
|
||||
}
|
||||
} else {
|
||||
expr.to_string()
|
||||
base.to_string()
|
||||
};
|
||||
|
||||
match step {
|
||||
Some(s) => format!("{}/{}", converted, s),
|
||||
None => converted,
|
||||
}
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user