Sync public tree to d0af060
This commit is contained in:
parent
8c11db913e
commit
a77e138a24
6 changed files with 185 additions and 11 deletions
|
|
@ -587,6 +587,7 @@ pub enum CoordinatorResponse {
|
|||
actor: WorkflowActor,
|
||||
placement: Placement,
|
||||
assignment: TaskAssignment,
|
||||
charged_spawns: u64,
|
||||
},
|
||||
TaskAssignment {
|
||||
assignment: Option<TaskAssignment>,
|
||||
|
|
@ -607,6 +608,7 @@ pub enum CoordinatorResponse {
|
|||
process: ProcessId,
|
||||
epoch: u64,
|
||||
actor: WorkflowActor,
|
||||
charged_spawns: u64,
|
||||
},
|
||||
NodeReconnected {
|
||||
node: NodeId,
|
||||
|
|
@ -761,6 +763,8 @@ pub struct CoordinatorService {
|
|||
download_meter: ResourceMeter,
|
||||
debug_limits: ResourceLimits,
|
||||
debug_meter: ResourceMeter,
|
||||
workflow_limits: ResourceLimits,
|
||||
workflow_meter: ResourceMeter,
|
||||
}
|
||||
|
||||
impl CoordinatorService {
|
||||
|
|
@ -790,6 +794,8 @@ impl CoordinatorService {
|
|||
download_meter: ResourceMeter::default(),
|
||||
debug_limits: ResourceLimits::community_tier_defaults(),
|
||||
debug_meter: ResourceMeter::default(),
|
||||
workflow_limits: ResourceLimits::community_tier_defaults(),
|
||||
workflow_meter: ResourceMeter::default(),
|
||||
}
|
||||
}
|
||||
|
||||
|
|
@ -1356,6 +1362,8 @@ impl CoordinatorService {
|
|||
)
|
||||
.into());
|
||||
}
|
||||
self.workflow_meter
|
||||
.can_charge(&self.workflow_limits, LimitKind::Spawn, 1)?;
|
||||
let request = PlacementRequest {
|
||||
tenant: tenant.clone(),
|
||||
project: project.clone(),
|
||||
|
|
@ -1374,6 +1382,9 @@ impl CoordinatorService {
|
|||
};
|
||||
let nodes = self.node_descriptors.values().cloned().collect::<Vec<_>>();
|
||||
let placement = DefaultScheduler.place(&nodes, &request)?;
|
||||
self.workflow_meter
|
||||
.charge(&self.workflow_limits, LimitKind::Spawn, 1)?;
|
||||
let charged_spawns = self.workflow_meter.used(&LimitKind::Spawn);
|
||||
let assignment = TaskAssignment {
|
||||
tenant: tenant.clone(),
|
||||
project: project.clone(),
|
||||
|
|
@ -1399,6 +1410,7 @@ impl CoordinatorService {
|
|||
actor,
|
||||
placement,
|
||||
assignment,
|
||||
charged_spawns,
|
||||
})
|
||||
}
|
||||
CoordinatorRequest::PollTaskAssignment {
|
||||
|
|
@ -1562,6 +1574,8 @@ impl CoordinatorService {
|
|||
.into());
|
||||
}
|
||||
}
|
||||
self.workflow_meter
|
||||
.charge(&self.workflow_limits, LimitKind::Spawn, 1)?;
|
||||
self.process_cancellations
|
||||
.remove(&process_control_key(&tenant, &project, &process));
|
||||
self.task_cancellations.retain(
|
||||
|
|
@ -1598,6 +1612,7 @@ impl CoordinatorService {
|
|||
process,
|
||||
epoch: self.coordinator.coordinator_epoch(),
|
||||
actor,
|
||||
charged_spawns: self.workflow_meter.used(&LimitKind::Spawn),
|
||||
})
|
||||
}
|
||||
CoordinatorRequest::ReconnectNode {
|
||||
|
|
@ -3411,6 +3426,144 @@ mod tests {
|
|||
assert!(policy_error.to_string().contains("policy denied placement"));
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn service_checks_spawn_quota_before_process_or_task_work_starts() {
|
||||
let mut service = CoordinatorService::new(7);
|
||||
service.workflow_limits = ResourceLimits {
|
||||
limits: BTreeMap::from([(LimitKind::Spawn, 2)]),
|
||||
};
|
||||
|
||||
service
|
||||
.handle_request(CoordinatorRequest::AttachNode {
|
||||
tenant: "tenant".to_owned(),
|
||||
project: "project".to_owned(),
|
||||
node: "worker-linux".to_owned(),
|
||||
public_key: "worker-linux-public-key".to_owned(),
|
||||
})
|
||||
.unwrap();
|
||||
service
|
||||
.handle_request(CoordinatorRequest::ReportNodeCapabilities {
|
||||
tenant: "tenant".to_owned(),
|
||||
project: "project".to_owned(),
|
||||
node: "worker-linux".to_owned(),
|
||||
capabilities: linux_capabilities(),
|
||||
cached_environment_digests: Vec::new(),
|
||||
dependency_cache_digests: Vec::new(),
|
||||
source_snapshots: Vec::new(),
|
||||
artifact_locations: Vec::new(),
|
||||
direct_connectivity: true,
|
||||
online: true,
|
||||
})
|
||||
.unwrap();
|
||||
|
||||
let CoordinatorResponse::ProcessStarted { charged_spawns, .. } = service
|
||||
.handle_request(CoordinatorRequest::StartProcess {
|
||||
tenant: "tenant".to_owned(),
|
||||
project: "project".to_owned(),
|
||||
actor_user: None,
|
||||
actor_agent: None,
|
||||
agent_public_key_fingerprint: None,
|
||||
process: "vp-quota".to_owned(),
|
||||
restart: false,
|
||||
})
|
||||
.unwrap()
|
||||
else {
|
||||
panic!("expected process start within spawn quota");
|
||||
};
|
||||
assert_eq!(charged_spawns, 1);
|
||||
|
||||
let CoordinatorResponse::TaskLaunched { charged_spawns, .. } = service
|
||||
.handle_request(CoordinatorRequest::LaunchTask {
|
||||
tenant: "tenant".to_owned(),
|
||||
project: "project".to_owned(),
|
||||
actor_user: Some("user".to_owned()),
|
||||
actor_agent: None,
|
||||
agent_public_key_fingerprint: None,
|
||||
process: "vp-quota".to_owned(),
|
||||
task: "compile-linux".to_owned(),
|
||||
environment: None,
|
||||
environment_digest: None,
|
||||
required_capabilities: vec![Capability::Command],
|
||||
dependency_cache: None,
|
||||
source_snapshot: None,
|
||||
required_artifacts: Vec::new(),
|
||||
quota_available: true,
|
||||
policy_allowed: true,
|
||||
command: "cargo".to_owned(),
|
||||
command_args: vec!["test".to_owned()],
|
||||
artifact_path: "/vfs/artifacts/dap-output.txt".to_owned(),
|
||||
})
|
||||
.unwrap()
|
||||
else {
|
||||
panic!("expected task launch within spawn quota");
|
||||
};
|
||||
assert_eq!(charged_spawns, 2);
|
||||
|
||||
let denied_task = service
|
||||
.handle_request(CoordinatorRequest::LaunchTask {
|
||||
tenant: "tenant".to_owned(),
|
||||
project: "project".to_owned(),
|
||||
actor_user: Some("user".to_owned()),
|
||||
actor_agent: None,
|
||||
agent_public_key_fingerprint: None,
|
||||
process: "vp-quota".to_owned(),
|
||||
task: "compile-linux-denied".to_owned(),
|
||||
environment: None,
|
||||
environment_digest: None,
|
||||
required_capabilities: vec![Capability::Command],
|
||||
dependency_cache: None,
|
||||
source_snapshot: None,
|
||||
required_artifacts: Vec::new(),
|
||||
quota_available: true,
|
||||
policy_allowed: true,
|
||||
command: "cargo".to_owned(),
|
||||
command_args: vec!["test".to_owned()],
|
||||
artifact_path: "/vfs/artifacts/denied.txt".to_owned(),
|
||||
})
|
||||
.unwrap_err();
|
||||
assert!(denied_task.to_string().contains("Spawn"));
|
||||
assert_eq!(service.workflow_meter.used(&LimitKind::Spawn), 2);
|
||||
|
||||
let denied_task_key = task_control_key(
|
||||
&TenantId::from("tenant"),
|
||||
&ProjectId::from("project"),
|
||||
&ProcessId::from("vp-quota"),
|
||||
&NodeId::from("worker-linux"),
|
||||
&TaskId::from("compile-linux-denied"),
|
||||
);
|
||||
assert!(!service.active_tasks.contains(&denied_task_key));
|
||||
assert!(!service.task_placements.contains_key(&denied_task_key));
|
||||
assert_eq!(
|
||||
service
|
||||
.task_assignments
|
||||
.get(&(
|
||||
TenantId::from("tenant"),
|
||||
ProjectId::from("project"),
|
||||
NodeId::from("worker-linux"),
|
||||
))
|
||||
.map(VecDeque::len),
|
||||
Some(1)
|
||||
);
|
||||
|
||||
let denied_process = service
|
||||
.handle_request(CoordinatorRequest::StartProcess {
|
||||
tenant: "tenant".to_owned(),
|
||||
project: "other-project".to_owned(),
|
||||
actor_user: None,
|
||||
actor_agent: None,
|
||||
agent_public_key_fingerprint: None,
|
||||
process: "vp-denied".to_owned(),
|
||||
restart: false,
|
||||
})
|
||||
.unwrap_err();
|
||||
assert!(denied_process.to_string().contains("Spawn"));
|
||||
assert!(service
|
||||
.coordinator
|
||||
.active_process(&ProcessId::from("vp-denied"))
|
||||
.is_none());
|
||||
assert_eq!(service.workflow_meter.used(&LimitKind::Spawn), 2);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn service_attaches_node_starts_process_and_records_scoped_task_event() {
|
||||
let mut service = CoordinatorService::new(7);
|
||||
|
|
@ -3449,7 +3602,8 @@ mod tests {
|
|||
public_key_fingerprint: None,
|
||||
authenticated_without_browser: false,
|
||||
scopes: vec!["project:read".to_owned(), "project:run".to_owned()],
|
||||
}
|
||||
},
|
||||
charged_spawns: 1,
|
||||
}
|
||||
);
|
||||
|
||||
|
|
@ -4296,7 +4450,8 @@ mod tests {
|
|||
public_key_fingerprint: None,
|
||||
authenticated_without_browser: false,
|
||||
scopes: vec!["project:read".to_owned(), "project:run".to_owned()],
|
||||
}
|
||||
},
|
||||
charged_spawns: 2,
|
||||
}
|
||||
);
|
||||
}
|
||||
|
|
|
|||
Loading…
Add table
Add a link
Reference in a new issue