Skip to content
Draft
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
75 changes: 73 additions & 2 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 @@ -14,7 +14,7 @@ hakari-package = "workspace-hack"

[workspace.package]
edition = "2024"
version = "0.41.0"
version = "0.42.0"

[workspace.dependencies]
anyhow = { version = "1.0.102", default-features = false }
Expand Down
16 changes: 10 additions & 6 deletions opsqueue/Cargo.toml
Original file line number Diff line number Diff line change
Expand Up @@ -7,12 +7,12 @@ repository = "https://github.com/channable/opsqueue"
license = "MIT"

[lib]
name="opsqueue"
path="src/lib.rs"
name = "opsqueue"
path = "src/lib.rs"

[[bin]]
name="opsqueue"
path="app/main.rs"
name = "opsqueue"
path = "app/main.rs"
required-features = ["server-logic"]

[dependencies]
Expand Down Expand Up @@ -72,6 +72,7 @@ humantime.workspace = true
dashmap.workspace = true
sqlformat.workspace = true
workspace-hack.workspace = true
async-trait = "0.1.91"

# See more keys and their definitions at https://doc.rust-lang.org/cargo/reference/manifest.html

Expand All @@ -80,6 +81,8 @@ workspace = true

[dev-dependencies]
insta.workspace = true
wiremock = "0.6.5"
tower = "0.5.3"

[[bench]]
name = "chunks_select"
Expand All @@ -98,8 +101,9 @@ server-logic = [
"dep:tower-http",
"dep:axum-prometheus",
"dep:sentry",
"dep:sentry-tracing"
]
"dep:sentry-tracing",
"dep:reqwest",
]
# Dependencies only in use by the client libraries:
client-logic = [
"dep:reqwest",
Expand Down
27 changes: 26 additions & 1 deletion opsqueue/app/main.rs
Original file line number Diff line number Diff line change
Expand Up @@ -4,12 +4,15 @@ use opentelemetry_otlp::SpanExporter;
use opentelemetry_resource_detectors::HostResourceDetector;
use opentelemetry_resource_detectors::{OsResourceDetector, ProcessResourceDetector};
use opentelemetry_sdk::trace::{RandomIdGenerator, Sampler, SdkTracerProvider};
use opsqueue::common::extension::{CoreApi, Extension};
use opsqueue::delegation::extension::DelegationExtension;
use opsqueue::tracing::as_dyn_error;
use opsqueue::{common::submission::db::periodically_cleanup_old, config::Config, prometheus};
use std::{
sync::{Arc, atomic::AtomicBool},
time::Duration,
};
use tokio::sync::Notify;
use tokio_util::sync::CancellationToken;
use tracing::level_filters::LevelFilter;

Expand Down Expand Up @@ -52,6 +55,24 @@ pub async fn async_main() {
.await
.expect("Timed out while initiating the database");

let notify_on_insert = Arc::new(Notify::new());
let (submission_status_changed_tx, _) = tokio::sync::broadcast::channel(128);

let core_api = CoreApi::new(
db_pool.clone(),
notify_on_insert.clone(),
submission_status_changed_tx.clone(),
);

let mut extensions: Vec<Box<dyn Extension>> = Vec::new();
if let Some(delegation_server_url) = &config.delegation_server_url {
extensions.push(Box::new(DelegationExtension::new(
cancellation_token.clone(),
core_api,
delegation_server_url.clone(),
)));
}

moro_local::async_scope!(|scope| {
let checkpoint_handle = scope.spawn(db_pool.periodically_checkpoint_wal());

Expand All @@ -63,10 +84,14 @@ pub async fn async_main() {
&cancellation_token,
&app_healthy_flag,
prometheus_config,
notify_on_insert,
submission_status_changed_tx,
&extensions,
));

let max_age = config.max_submission_age.into();
let cleanup_handle = scope.spawn(periodically_cleanup_old(db_pool.writer_pool(), max_age));

let cleanup_handle = scope.spawn(periodically_cleanup_old(&db_pool, max_age, &extensions));

let prometheus_handle = scope.spawn(prometheus::periodically_calculate_scaling_metrics(
&db_pool,
Expand Down
Original file line number Diff line number Diff line change
@@ -0,0 +1 @@
DROP TABLE submissions_external_task;
Original file line number Diff line number Diff line change
@@ -0,0 +1,6 @@
CREATE TABLE submissions_external_task
(
submission_id BIGINT NOT NULL UNIQUE,
task_id TEXT NOT NULL UNIQUE,
last_status_sent TEXT
);
Binary file modified opsqueue/opsqueue_example_database_schema.db
Binary file not shown.
Loading
Loading