Skip to content

Commit f0e4026

Browse files
committed
[ACTP] orchestrate par-control tasks
1 parent b79c0d2 commit f0e4026

7 files changed

Lines changed: 1533 additions & 7 deletions

File tree

pkg/privateactionrunner/par-control/BUILD.bazel

Lines changed: 2 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -31,6 +31,7 @@ rust_library(
3131
"src/jwt.rs",
3232
"src/lib.rs",
3333
"src/opms.rs",
34+
"src/orchestrator.rs",
3435
"src/platform.rs",
3536
"src/procmgr.rs",
3637
"src/proto.rs",
@@ -126,6 +127,7 @@ rust_test(
126127
target_compatible_with = _LINUX_OR_WINDOWS,
127128
deps = [
128129
"@crates//:tempfile",
130+
"@crates//:tokio",
129131
"@crates//:tokio-stream",
130132
],
131133
)

pkg/privateactionrunner/par-control/Cargo.toml

Lines changed: 3 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -25,6 +25,7 @@ dd-procmgr-client.workspace = true
2525
hyper-util.workspace = true
2626
log.workspace = true
2727
native-tls.workspace = true
28+
# Used to convert the Agent IPC key from SEC1 to PKCS8 for native-tls.
2829
openssl.workspace = true
2930
p256 = { workspace = true, features = ["ecdsa", "pem", "pkcs8", "std"] }
3031
prost.workspace = true
@@ -37,6 +38,8 @@ tokio = { workspace = true, features = [
3738
"net",
3839
"rt-multi-thread",
3940
"signal",
41+
"sync",
42+
"time",
4043
] }
4144
tokio-native-tls.workspace = true
4245
tokio-stream = { workspace = true, features = ["net"] }

pkg/privateactionrunner/par-control/src/bins/par-control.rs

Lines changed: 58 additions & 7 deletions
Original file line numberDiff line numberDiff line change
@@ -7,9 +7,15 @@ use anyhow::Result;
77
use clap::Parser;
88
use par_control::bootstrap;
99
use par_control::config::Launch;
10+
use par_control::executor::ExecutorDispatcher;
11+
use par_control::jwt::{Es256Signer, JwtSigner};
12+
use par_control::opms::{HttpOpms, HttpOpmsConfig};
13+
use par_control::orchestrator::{Orchestrator, Params};
1014
use par_control::platform;
15+
use par_control::procmgr::ProcmgrLifecycle;
1116
use std::path::PathBuf;
1217
use std::process::ExitCode;
18+
use std::sync::Arc;
1319

1420
#[derive(Parser)]
1521
#[command(name = "par-control", about = "Private Action Runner control plane")]
@@ -42,31 +48,76 @@ async fn main() -> ExitCode {
4248
async fn run() -> Result<()> {
4349
let cli = Cli::parse();
4450
// Resolve the log level first so a config error is logged at that level.
51+
// A logging failure does not prevent the runner from starting.
4552
let launch = Launch::from_yaml_file(&cli.config);
4653
let log_level = launch
4754
.as_ref()
4855
.map_or(log::LevelFilter::Info, |launch| launch.log_level);
49-
if let Err(error) = dd_agent_log::init(dd_agent_log::LogConfig {
56+
if let Err(e) = dd_agent_log::init(dd_agent_log::LogConfig {
5057
logger_name: "PAR-CONTROL",
5158
level: log_level.to_level().unwrap_or(log::Level::Error),
5259
// dd-procmgrd redirects stdout and stderr per its process definition.
5360
log_file: None,
5461
}) {
55-
eprintln!("par-control: could not initialize the logger: {error}");
62+
eprintln!("par-control: could not initialize the logger: {e}");
5663
}
5764
log::set_max_level(log_level);
5865

5966
if !launch?.gate.split_mode {
60-
log::info!("private_action_runner split mode is disabled; par-control is exiting");
67+
log::info!(
68+
"private_action_runner.split_enabled is not enabled; \
69+
the monolithic runner owns OPMS polling. Exiting."
70+
);
6171
return Ok(());
6272
}
6373

64-
let _config =
74+
let config =
6575
bootstrap::load_config_with_bootstrap(&cli.config, &cli.ensure_enrollment_command)?;
6676

67-
log::info!("par-control started");
68-
shutdown_signal().await;
69-
log::info!("par-control is exiting");
77+
let signer: Arc<dyn JwtSigner> = Arc::new(Es256Signer::new(
78+
config.identity.org_id,
79+
config.identity.runner_id.clone(),
80+
&config.identity.private_key,
81+
)?);
82+
83+
let opms = Arc::new(HttpOpms::new(
84+
config.opms_base_url.clone(),
85+
signer,
86+
HttpOpmsConfig {
87+
runner_version: config.runner_version.clone(),
88+
modes: config.modes.clone(),
89+
timeout: config.opms_request_timeout,
90+
proxy: config.proxy.clone(),
91+
tls: config.tls.clone(),
92+
extra_headers: config.opms_extra_headers.clone(),
93+
},
94+
)?);
95+
let lifecycle = Arc::new(ProcmgrLifecycle::new(
96+
&config.procmgr_socket,
97+
config.executor_process_name.clone(),
98+
));
99+
// The IPC certificate is loaded lazily because the executor may create it.
100+
let dispatcher = Arc::new(ExecutorDispatcher::new(
101+
&config.executor_socket,
102+
Some(&config.ipc_cert_file),
103+
));
104+
105+
let params = Params::from_config(&config);
106+
let orchestrator = Orchestrator::new(opms, lifecycle, dispatcher, params);
107+
108+
log::info!(
109+
"par-control starting: version={} urn={} opms={} executor_socket={} procmgr_socket={} ipc_cert={}",
110+
config.runner_version,
111+
config.identity.urn,
112+
config.opms_base_url,
113+
config.executor_socket.display(),
114+
config.procmgr_socket.display(),
115+
config.ipc_cert_file.display(),
116+
);
117+
118+
orchestrator.run(shutdown_signal()).await;
119+
log::info!("par-control stopped");
120+
log::logger().flush();
70121
Ok(())
71122
}
72123

pkg/privateactionrunner/par-control/src/lib.rs

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -9,6 +9,7 @@ pub mod executor;
99
pub mod identity;
1010
pub mod jwt;
1111
pub mod opms;
12+
pub mod orchestrator;
1213
pub mod platform;
1314
pub mod procmgr;
1415
pub mod proto;

0 commit comments

Comments
 (0)