feat: dvfs agent interface
This commit is contained in:
@@ -0,0 +1,68 @@
|
||||
package nextflow.k8s
|
||||
|
||||
import groovy.util.logging.Slf4j
|
||||
|
||||
import java.net.http.HttpClient
|
||||
import java.net.http.HttpRequest
|
||||
import java.net.http.HttpResponse
|
||||
|
||||
/**
|
||||
* The DVFS client uses the deployed DVFS agents to control the operating frequency of the worker nodes.
|
||||
*/
|
||||
@Slf4j
|
||||
class K8sDVFSClient {
|
||||
private HttpClient httpClient
|
||||
|
||||
K8sDVFSClient() {
|
||||
this.httpClient = HttpClient.newBuilder().build()
|
||||
}
|
||||
|
||||
OptionalInt getNodeCurrentFrequency(String node) {
|
||||
log.debug("Getting current frequency of node ${node}")
|
||||
return getNodeFrequency(node, "current")
|
||||
}
|
||||
|
||||
OptionalInt getNodeMaxFrequency(String node) {
|
||||
log.debug("Getting max frequency of node ${node}")
|
||||
return getNodeFrequency(node, "max")
|
||||
}
|
||||
|
||||
OptionalInt getNodeMinFrequency(String node) {
|
||||
log.debug("Getting min frequency of node ${node}")
|
||||
return getNodeFrequency(node, "min")
|
||||
}
|
||||
|
||||
private OptionalInt getNodeFrequency(String node, String endpoint) {
|
||||
HttpRequest request = HttpRequest.newBuilder()
|
||||
.uri(new URI("http://${nodeName}/cpu/frequency/${endpoint}"))
|
||||
.GET()
|
||||
.build()
|
||||
HttpResponse<String> response = httpClient.send(request, HttpResponse.BodyHandlers.ofString())
|
||||
if (response.statusCode() != 200) {
|
||||
log.error("Request GET ${request.uri().toString()} returned ${response.statusCode()}")
|
||||
return OptionalInt.empty()
|
||||
}
|
||||
try {
|
||||
double parsed = Double.parseDouble(response.body())
|
||||
return OptionalInt.of(parsed)
|
||||
} catch (NumberFormatException ex) {
|
||||
log.error("Unexpected response ${response.body()} - ${ex.message}")
|
||||
return OptionalInt.empty()
|
||||
}
|
||||
}
|
||||
|
||||
boolean setNodeFrequency(String node, int frequency) {
|
||||
String body = String.valueOf(frequency)
|
||||
|
||||
HttpRequest request = HttpRequest.newBuilder()
|
||||
.uri(new URI("http://${nodeName}/cpu/frequency/current"))
|
||||
.PUT(HttpRequ3est.BodyPublishers.ofString(body))
|
||||
.build()
|
||||
HttpResponse<Void> response = httpClient.send(request, HttpResponse.BodyHandlers.discarding())
|
||||
if (response.statusCode() != 200) {
|
||||
log.error("Request PUT ${request.uri().toString()} returned ${response.statusCode()}")
|
||||
return false
|
||||
}
|
||||
return true
|
||||
}
|
||||
}
|
||||
Reference in New Issue
Block a user