diff --git a/nextflow/plugins/nf-k8s/src/main/nextflow/k8s/K8sDVFSClient.groovy b/nextflow/plugins/nf-k8s/src/main/nextflow/k8s/K8sDVFSClient.groovy new file mode 100644 index 0000000..b0a6116 --- /dev/null +++ b/nextflow/plugins/nf-k8s/src/main/nextflow/k8s/K8sDVFSClient.groovy @@ -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 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 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 + } +}