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
458 changes: 331 additions & 127 deletions Cargo.lock

Large diffs are not rendered by default.

1 change: 1 addition & 0 deletions Cargo.toml
Original file line number Diff line number Diff line change
Expand Up @@ -10,6 +10,7 @@ fail-on-warnings = []
[dependencies]
es-entity = "0.9.0"
sqlx-ledger = { version = "0.11.5", features = ["otel"] }
job_crate = { package = "job", version = "0.2.1" }

anyhow = "1.0.82"
bitcoincore-rpc = "0.17.0"
Expand Down
61 changes: 61 additions & 0 deletions migrations/20250904065521_job_setup.sql
Original file line number Diff line number Diff line change
@@ -0,0 +1,61 @@
CREATE TABLE jobs (
id UUID PRIMARY KEY,
unique_per_type BOOLEAN NOT NULL,
job_type VARCHAR NOT NULL,
created_at TIMESTAMPTZ NOT NULL DEFAULT NOW()
);
CREATE UNIQUE INDEX idx_unique_job_type ON jobs (job_type) WHERE unique_per_type = TRUE;

CREATE TABLE job_events (
id UUID NOT NULL REFERENCES jobs(id),
sequence INT NOT NULL,
event_type VARCHAR NOT NULL,
event JSONB NOT NULL,
context JSONB DEFAULT NULL,
recorded_at TIMESTAMPTZ NOT NULL DEFAULT NOW(),
UNIQUE(id, sequence)
);

CREATE TYPE JobExecutionState AS ENUM ('pending', 'running');

CREATE TABLE job_executions (
id UUID REFERENCES jobs(id) NOT NULL UNIQUE,
job_type VARCHAR NOT NULL,
poller_instance_id UUID,
attempt_index INT NOT NULL DEFAULT 1,
state JobExecutionState NOT NULL DEFAULT 'pending',
execution_state_json JSONB,
execute_at TIMESTAMPTZ,
alive_at TIMESTAMPTZ NOT NULL,
created_at TIMESTAMPTZ NOT NULL
);

CREATE INDEX idx_job_executions_poller_instance
ON job_executions(poller_instance_id)
WHERE state = 'running';

CREATE OR REPLACE FUNCTION notify_job_execution_insert() RETURNS TRIGGER AS $$
BEGIN
PERFORM pg_notify('job_execution', '');
RETURN NULL;
END;
$$ LANGUAGE plpgsql;

CREATE OR REPLACE FUNCTION notify_job_execution_update() RETURNS TRIGGER AS $$
BEGIN
IF NEW.execute_at IS DISTINCT FROM OLD.execute_at THEN
PERFORM pg_notify('job_execution', '');
END IF;
RETURN NULL;
END;
$$ LANGUAGE plpgsql;

CREATE TRIGGER job_executions_notify_insert_trigger
AFTER INSERT ON job_executions
FOR EACH STATEMENT
EXECUTE FUNCTION notify_job_execution_insert();

CREATE TRIGGER job_executions_notify_update_trigger
AFTER UPDATE ON job_executions
FOR EACH STATEMENT
EXECUTE FUNCTION notify_job_execution_update();
22 changes: 18 additions & 4 deletions src/admin/app.rs
Original file line number Diff line number Diff line change
@@ -1,7 +1,9 @@
use tracing::instrument;

use super::{error::*, keys::*};
use crate::{account::*, dev_constants, ledger::Ledger, primitives::bitcoin, profile::*};
use crate::{
account::*, dev_constants, job_svc::JobSvc, ledger::Ledger, primitives::bitcoin, profile::*,
};

const BOOTSTRAP_KEY_NAME: &str = "admin_bootstrap_key";

Expand All @@ -11,16 +13,18 @@ pub struct AdminApp {
profiles: Profiles,
ledger: Ledger,
network: bitcoin::Network,
job_svc: JobSvc,
}

impl AdminApp {
pub fn new(pool: sqlx::PgPool, network: bitcoin::Network) -> Self {
pub fn new(pool: sqlx::PgPool, network: bitcoin::Network, job_svc: JobSvc) -> Self {
Self {
keys: AdminApiKeys::new(&pool),
accounts: Accounts::new(&pool),
profiles: Profiles::new(&pool),
ledger: Ledger::new(&pool),
network,
job_svc,
}
}
}
Expand All @@ -43,14 +47,19 @@ impl AdminApp {
.await?;
let new_profile = NewProfile::builder()
.account_id(account.id)
.name(account.name)
.name(account.name.clone())
.build()
.expect("Couldn't build NewProfile");
let profile = self.profiles.create_in_op(&mut op, new_profile).await?;
let profile_key = self
.profiles
.create_key_for_profile_in_op(&mut op, profile, true)
.await?;

self.job_svc
.spawn_outbox_handler_in_op(&mut op, account.id, account.journal_id())
.await?;

op.commit().await?;
Ok((admin_key, profile_key))
}
Expand Down Expand Up @@ -81,14 +90,19 @@ impl AdminApp {
.await?;
let new_profile = NewProfile::builder()
.account_id(account.id)
.name(account.name)
.name(account.name.clone())
.build()
.expect("Couldn't build NewProfile");
let profile = self.profiles.create_in_op(&mut op, new_profile).await?;
let key = self
.profiles
.create_key_for_profile_in_op(&mut op, profile, false)
.await?;

self.job_svc
.spawn_outbox_handler_in_op(&mut op, account.id, account.journal_id())
.await?;

op.commit().await?;
Ok(key)
}
Expand Down
6 changes: 4 additions & 2 deletions src/admin/error.rs
Original file line number Diff line number Diff line change
@@ -1,8 +1,8 @@
use thiserror::Error;

use crate::{
account::error::AccountError, app::error::ApplicationError, ledger::error::LedgerError,
profile::error::ProfileError,
account::error::AccountError, app::error::ApplicationError, job_svc::JobSvcError,
ledger::error::LedgerError, profile::error::ProfileError,
};

#[allow(clippy::large_enum_variant)]
Expand All @@ -22,6 +22,8 @@ pub enum AdminApiError {
ProfileError(#[from] ProfileError),
#[error("{0}")]
LedgerError(#[from] LedgerError),
#[error("{0}")]
JobSvcError(#[from] JobSvcError),
#[error("AdminApiError - DevBootstrapError: {0}")]
DevBootstrapError(#[from] anyhow::Error),
}
Expand Down
8 changes: 5 additions & 3 deletions src/admin/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -4,7 +4,7 @@ pub mod error;
mod keys;
mod server;

use crate::{dev_constants, primitives::bitcoin, token_store};
use crate::{dev_constants, job_svc::JobSvc, primitives::bitcoin, token_store};

pub use app::*;
pub use config::*;
Expand All @@ -17,8 +17,9 @@ pub async fn run_dev(
config: AdminApiConfig,
network: bitcoin::Network,
bria_home: String,
job_svc: JobSvc,
) -> Result<(), AdminApiError> {
let app = AdminApp::new(pool, network);
let app = AdminApp::new(pool, network, job_svc);
let (admin_key, profile_key) = app.dev_bootstrap().await?;
token_store::store_admin_token(&bria_home, &admin_key.key)?;
println!("Admin API key");
Expand All @@ -44,8 +45,9 @@ pub async fn run(
pool: sqlx::PgPool,
config: AdminApiConfig,
network: bitcoin::Network,
job_svc: JobSvc,
) -> Result<(), AdminApiError> {
let app = AdminApp::new(pool, network);
let app = AdminApp::new(pool, network, job_svc);
server::start(config, app).await?;
Ok(())
}
3 changes: 3 additions & 0 deletions src/app/error.rs
Original file line number Diff line number Diff line change
Expand Up @@ -9,6 +9,7 @@ use crate::{
descriptor::error::DescriptorError,
fees::error::FeeEstimationError,
job::error::JobError,
job_svc::JobSvcError,
ledger::error::LedgerError,
outbox::error::OutboxError,
payout::error::PayoutError,
Expand Down Expand Up @@ -44,6 +45,8 @@ pub enum ApplicationError {
#[error("{0}")]
JobError(#[from] JobError),
#[error("{0}")]
JobSvcError(#[from] JobSvcError),
#[error("{0}")]
OutboxError(#[from] OutboxError),
#[error("{0}")]
UtxoError(#[from] UtxoError),
Expand Down
39 changes: 11 additions & 28 deletions src/app/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -17,6 +17,7 @@ use crate::{
descriptor::*,
fees::{self, *},
job,
job_svc::*,
ledger::*,
outbox::*,
payout::*,
Expand All @@ -32,6 +33,7 @@ use crate::{
#[allow(dead_code)]
pub struct App {
_runner: JobRunnerHandle,
job_svc: JobSvc,
outbox: Outbox,
profiles: Profiles,
xpubs: XPubs,
Expand Down Expand Up @@ -68,6 +70,9 @@ impl App {
)
.await?;
let fees_client = FeesClient::new(config.fees.clone());

let job_svc = JobSvc::init(pool.clone(), outbox.clone(), ledger.clone()).await?;

let runner = job::start_job_runner(
&pool,
outbox.clone(),
Expand All @@ -92,12 +97,9 @@ impl App {
config.jobs.process_all_payout_queues_delay,
)
.await?;
Self::spawn_respawn_all_outbox_handlers(
pool.clone(),
config.jobs.respawn_all_outbox_handlers_delay,
)
.await?;

let app = Self {
job_svc,
outbox,
profiles: Profiles::new(&pool),
xpubs,
Expand Down Expand Up @@ -127,6 +129,10 @@ impl App {
self.config.blockchain.network
}

pub fn job_svc(&self) -> &JobSvc {
&self.job_svc
}

#[instrument(name = "app.authenticate", skip_all, err)]
pub async fn authenticate(&self, key: &str) -> Result<Profile, ApplicationError> {
let profile = self.profiles.find_by_key(key).await?;
Expand Down Expand Up @@ -1080,27 +1086,4 @@ impl App {
});
Ok(())
}

#[instrument(
name = "app.spawn_respawn_all_outbox_handlers",
level = "trace",
skip_all,
err
)]
async fn spawn_respawn_all_outbox_handlers(
pool: sqlx::PgPool,
delay: std::time::Duration,
) -> Result<(), ApplicationError> {
tokio::spawn(async move {
loop {
let _ = job::spawn_respawn_all_outbox_handlers(
&pool,
std::time::Duration::from_secs(1),
)
.await;
tokio::time::sleep(delay).await;
}
});
Ok(())
}
}
9 changes: 7 additions & 2 deletions src/cli/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -1071,20 +1071,25 @@ async fn run_cmd(
let mut handles = Vec::new();
let pool = init_pool(&db).await?;

let main_app = crate::app::App::run(pool.clone(), app.clone()).await?;

let admin_job_svc = main_app.job_svc().clone();

let admin_send = send.clone();
let admin_pool = pool.clone();
let network = app.blockchain.network;
handles.push(tokio::spawn(async move {
let _ = admin_send.try_send(if dev {
super::admin::run_dev(admin_pool, admin, network, bria_home)
super::admin::run_dev(admin_pool, admin, network, bria_home, admin_job_svc)
.await
.context("Admin server error")
} else {
super::admin::run(admin_pool, admin, network)
super::admin::run(admin_pool, admin, network, admin_job_svc)
.await
.context("Admin server error")
});
}));

let api_send = send.clone();
handles.push(tokio::spawn(async move {
let _ = api_send.try_send(if dev {
Expand Down
Loading
Loading