Skip to content

Commit 6b7f4a8

Browse files
authored
Merge pull request #258 from trydirect/feature/agent-sweep
Feature/agent sweep
2 parents b04fbe0 + 0515e96 commit 6b7f4a8

6 files changed

Lines changed: 465 additions & 0 deletions

File tree

‎src/cli/stacker_client.rs‎

Lines changed: 49 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -4096,6 +4096,21 @@ pub fn build_project_body(config: &StackerConfig) -> serde_json::Value {
40964096
service_apps.push(service_to_app_json(svc, &network_ids));
40974097
}
40984098

4099+
// `project_app.config_contract` is persisted by a separate server-side
4100+
// accessor. Include the full contract on each generated app so `stacker
4101+
// sync` does not silently discard the policy from stacker.yml.
4102+
if let Ok(config_contract) = serde_json::to_value(&config.config_contract) {
4103+
let has_services = config_contract
4104+
.get("services")
4105+
.and_then(serde_json::Value::as_object)
4106+
.is_some_and(|services| !services.is_empty());
4107+
if has_services {
4108+
for app in web_apps.iter_mut().chain(service_apps.iter_mut()) {
4109+
app["config_contract"] = config_contract.clone();
4110+
}
4111+
}
4112+
}
4113+
40994114
serde_json::json!({
41004115
"custom": {
41014116
"custom_stack_code": stack_code,
@@ -5129,6 +5144,40 @@ mod tests {
51295144
);
51305145
}
51315146

5147+
#[test]
5148+
fn build_project_body_includes_config_contract_on_apps() {
5149+
let mut config = crate::cli::config_parser::ConfigBuilder::new()
5150+
.name("contract-project")
5151+
.app_image("nginx:1.27")
5152+
.build()
5153+
.expect("config should build");
5154+
config.config_contract = serde_json::from_value(serde_json::json!({
5155+
"services": {
5156+
"app": {
5157+
"fields": {
5158+
"JWT_SECRET": {
5159+
"mutability": "generated",
5160+
"type": "hex",
5161+
"length": 32
5162+
}
5163+
}
5164+
}
5165+
}
5166+
}))
5167+
.expect("config contract should deserialize");
5168+
5169+
let body = build_project_body(&config);
5170+
let contract = &body["custom"]["web"][0]["config_contract"];
5171+
assert_eq!(
5172+
contract["services"]["app"]["fields"]["JWT_SECRET"]["mutability"],
5173+
"generated"
5174+
);
5175+
assert_eq!(
5176+
contract["services"]["app"]["fields"]["JWT_SECRET"]["type"],
5177+
"hex"
5178+
);
5179+
}
5180+
51325181
#[test]
51335182
fn pipe_instance_request_serializes_adapter_references() {
51345183
let request = CreatePipeInstanceApiRequest {

‎src/db/agent.rs‎

Lines changed: 69 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -239,6 +239,75 @@ pub async fn delete(pool: &PgPool, agent_id: Uuid) -> Result<(), String> {
239239
})
240240
}
241241

242+
/// Delete agents whose deployment is gone (or soft-deleted) and that show no
243+
/// sign of life within the retention window.
244+
///
245+
/// "No sign of life" means both:
246+
/// - `last_heartbeat` is NULL or older than the window, AND
247+
/// - no `audit_log` row references the agent (by id or deployment_hash)
248+
/// within the same window.
249+
///
250+
/// The second condition protects agents that are alive but failing
251+
/// authentication — `last_heartbeat` only advances on successful `wait`/`report`,
252+
/// while `audit_log` captures `auth_failure` entries.
253+
#[tracing::instrument(name = "Sweep dead agents", skip(pool))]
254+
pub async fn sweep_dead(pool: &PgPool, retention_days: i32) -> Result<u64, String> {
255+
let result = sqlx::query(
256+
r#"
257+
DELETE FROM agents a
258+
WHERE NOT EXISTS (
259+
SELECT 1 FROM deployment d
260+
WHERE d.deployment_hash = a.deployment_hash
261+
AND d.deleted IS NOT TRUE
262+
)
263+
AND (a.last_heartbeat IS NULL
264+
OR a.last_heartbeat < NOW() - make_interval(days => $1))
265+
AND NOT EXISTS (
266+
SELECT 1 FROM audit_log l
267+
WHERE (l.agent_id = a.id OR l.deployment_hash = a.deployment_hash)
268+
AND l.created_at > NOW() - make_interval(days => $1)
269+
)
270+
"#,
271+
)
272+
.bind(retention_days)
273+
.execute(pool)
274+
.await
275+
.map_err(|err| {
276+
tracing::error!("Failed to sweep dead agents: {:?}", err);
277+
format!("Database error: {}", err)
278+
})?;
279+
280+
Ok(result.rows_affected())
281+
}
282+
283+
/// Delete agents whose `deployment_hash` is structurally invalid.
284+
///
285+
/// A valid hash matches `deployment_<uuid>` (36-char UUID with hyphens).
286+
/// Rows with an invalid hash can never be matched by their own agent — the
287+
/// lookup in `fetch_by_deployment_hash` will never find them. Among the 18
288+
/// currently broken rows, 17 store a raw agent token (86-char base64url)
289+
/// instead of a hash, leaking the secret in plaintext.
290+
///
291+
/// Returns its own count separate from `sweep_dead` so the caller can log
292+
/// exactly how many malformed rows were removed.
293+
#[tracing::instrument(name = "Sweep malformed agent rows", skip(pool))]
294+
pub async fn sweep_malformed(pool: &PgPool) -> Result<u64, String> {
295+
let result = sqlx::query(
296+
r#"
297+
DELETE FROM agents
298+
WHERE deployment_hash !~ '^deployment_[0-9a-fA-F]{8}-[0-9a-fA-F]{4}-[0-9a-fA-F]{4}-[0-9a-fA-F]{4}-[0-9a-fA-F]{12}$'
299+
"#,
300+
)
301+
.execute(pool)
302+
.await
303+
.map_err(|err| {
304+
tracing::error!("Failed to sweep malformed agents: {:?}", err);
305+
format!("Database error: {}", err)
306+
})?;
307+
308+
Ok(result.rows_affected())
309+
}
310+
242311
pub async fn log_audit(
243312
pool: &PgPool,
244313
audit_log: models::AuditLog,

‎src/services/agent_sweeper.rs‎

Lines changed: 76 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,76 @@
1+
//! Housekeeping for dead and malformed agent rows.
2+
//!
3+
//! The `agents` table accumulates rows that no longer serve a purpose: agents
4+
//! whose deployment was deleted or soft-deleted, agents that never sent a
5+
//! heartbeat, and rows whose `deployment_hash` is structurally invalid (e.g. a
6+
//! raw agent token stored instead of a hash). This sweeper removes them on a
7+
//! daily cadence.
8+
//!
9+
//! A row is considered "dead" only when **both** of the following hold:
10+
//!
11+
//! 1. The deployment is gone — either no matching row in `deployment`, or the
12+
//! row exists with `deleted IS TRUE`.
13+
//! 2. There is no sign of life within the retention window — neither a
14+
//! `last_heartbeat` nor an `audit_log` entry (the latter protects agents
15+
//! that are alive but failing authentication, since `auth_failure` is
16+
//! recorded in `audit_log` while `last_heartbeat` only advances on
17+
//! successful `wait`/`report`).
18+
//!
19+
//! Malformed rows (invalid `deployment_hash`) are deleted unconditionally,
20+
//! without a retention window — such an agent can never authenticate by its own
21+
//! hash.
22+
//!
23+
//! Agents whose deployment is alive but silent are **not** touched: the row is
24+
//! the agent's identity, and removing it would require a full reinstall.
25+
26+
use std::time::Duration;
27+
28+
use sqlx::PgPool;
29+
30+
use crate::db;
31+
32+
/// How often to sweep. Rows appear rarely; daily is ample.
33+
const TICK: Duration = Duration::from_secs(86_400);
34+
35+
/// How long a dead agent row is kept before removal. Long enough that a
36+
/// temporarily stopped server can come back without losing its identity.
37+
const RETENTION_DAYS: i32 = 30;
38+
39+
pub fn spawn(pg_pool: PgPool) {
40+
tokio::spawn(async move {
41+
tracing::info!(
42+
"agent_sweeper started (tick={:?}, retention={} days)",
43+
TICK,
44+
RETENTION_DAYS
45+
);
46+
loop {
47+
// Sleep first: startup is busy enough, and nothing here is urgent.
48+
tokio::time::sleep(TICK).await;
49+
50+
// Malformed rows are independent of retention — log at warn because
51+
// new appearances indicate a write path that still needs fixing.
52+
match db::agent::sweep_malformed(&pg_pool).await {
53+
Ok(0) => {}
54+
Ok(count) => tracing::warn!(
55+
"agent_sweeper: removed {} malformed agent row(s) — \
56+
the write path producing these has not been found yet",
57+
count
58+
),
59+
Err(err) => {
60+
tracing::warn!("agent_sweeper: malformed sweep error: {}", err)
61+
}
62+
}
63+
64+
match db::agent::sweep_dead(&pg_pool, RETENTION_DAYS).await {
65+
Ok(0) => tracing::debug!("agent_sweeper: nothing to remove"),
66+
Ok(count) => tracing::info!(
67+
"agent_sweeper: removed {} dead agent row(s)",
68+
count
69+
),
70+
Err(err) => {
71+
tracing::warn!("agent_sweeper: dead sweep error: {}", err)
72+
}
73+
}
74+
}
75+
});
76+
}

‎src/services/mod.rs‎

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -1,4 +1,5 @@
11
pub mod agent_dispatcher;
2+
pub mod agent_sweeper;
23
pub mod agent_token;
34
pub mod config_renderer;
45
pub mod dag_executor;

‎src/startup.rs‎

Lines changed: 5 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -122,6 +122,11 @@ pub async fn run(
122122
// it touches nothing the dashboard shows, and skipping it costs a slowly
123123
// growing table rather than correctness.
124124
crate::services::deployment_container_sweeper::spawn(api_pool.get_ref().clone());
125+
// Removes dead agent rows (no live deployment, no recent activity) and
126+
// rows with structurally invalid deployment_hash. Runs daily; skipping it
127+
// means the agents table slowly accumulates noise that obscures real
128+
// problems.
129+
crate::services::agent_sweeper::spawn(api_pool.get_ref().clone());
125130

126131
let payout_provider = crate::services::init_payout_provider(&settings.payouts)
127132
.map_err(|err| std::io::Error::new(std::io::ErrorKind::InvalidInput, err.to_string()))?;

0 commit comments

Comments
 (0)