Compilation Parts: fast-check for already uploaded files

This commit is contained in:
Vladislav Rassokhin
2018-11-21 16:30:09 +03:00
parent 9e094569d9
commit fd6c866e81
2 changed files with 81 additions and 22 deletions
@@ -1,15 +1,17 @@
// Copyright 2000-2018 JetBrains s.r.o. Use of this source code is governed by the Apache 2.0 license that can be found in the LICENSE file.
package org.jetbrains.intellij.build.impl
import com.google.gson.Gson
import com.intellij.openapi.util.io.StreamUtil
import com.intellij.openapi.util.text.StringUtil
import groovy.transform.CompileStatic
import org.apache.http.client.methods.CloseableHttpResponse
import org.apache.http.client.methods.HttpGet
import org.apache.http.client.methods.HttpPost
import org.apache.http.client.methods.HttpPut
import org.apache.http.entity.ContentType
import org.apache.http.entity.FileEntity
import org.apache.http.entity.StringEntity
import org.apache.http.impl.client.CloseableHttpClient
import org.apache.http.impl.client.HttpClientBuilder
import org.apache.http.impl.client.LaxRedirectStrategy
@@ -66,22 +68,22 @@ class CompilationPartsUploader implements Closeable {
myMessages.error(message)
}
boolean upload(@NotNull final String path, @NotNull final File file) throws UploadException {
boolean upload(@NotNull final String path, @NotNull final File file, boolean sendHead) throws UploadException {
debug("Preparing to upload " + file + " to " + myServerUrl)
if (!file.exists()) {
throw new UploadException("The file " + file.getPath() + " does not exist")
}
int code = doHead(path)
if (code == 200) {
log("File '$path' already exist on server, nothing to upload")
return false
if (sendHead) {
int code = doHead(path)
if (code == 200) {
log("File '$path' already exist on server, nothing to upload")
return false
}
if (code != 404) {
error("HEAD $path responded with unexpected $code")
}
}
if (code != 404) {
error("HEAD $path responded with unexpected $code")
}
final String response = doPut(path, file)
if (StringUtil.isEmptyOrSpaces(response)) {
log("Performed '$path' upload.")
@@ -92,6 +94,43 @@ class CompilationPartsUploader implements Closeable {
return true
}
@CompileStatic
static class CheckFilesResponse {
List<String> found
List<String> missing
CheckFilesResponse() {
}
}
CheckFilesResponse getFoundAndMissingFiles(String metadataJson) {
String path = '/check-files'
CloseableHttpResponse response = null
try {
String url = myServerUrl + StringUtil.trimStart(path, '/')
debug("POST " + url)
def request = new HttpPost(url)
request.setEntity(new StringEntity(metadataJson, ContentType.APPLICATION_JSON))
response = myHttpClient.execute(request)
debug("POST code: ${response.getStatusLine().getStatusCode()}")
def responseString = EntityUtils.toString(response.getEntity(), ContentType.APPLICATION_JSON.charset)
def parsedResponse = new Gson().fromJson(responseString, CheckFilesResponse.class)
return parsedResponse
}
catch (Exception e) {
myMessages.warning("Failed to check for found and mising files ('$path'): ${e.message}")
return null
}
finally {
StreamUtil.closeStream(response)
}
}
@NotNull
private int doHead(String path) throws UploadException {
CloseableHttpResponse response = null
@@ -179,6 +179,14 @@ class CompilationPartsUtil {
}
}
// Prepare metadata for writing into file
CompilationPartsMetadata m = new CompilationPartsMetadata()
m.serverUrl = serverUrl
m.branch = branch
m.prefix = uploadPrefix
m.files = new TreeMap<String, String>(hashes)
String metadataJson = new Gson().toJson(m)
messages.block("Uploading archives") {
AtomicInteger uploadedCount = new AtomicInteger()
AtomicLong uploadedBytes = new AtomicLong()
@@ -188,18 +196,37 @@ class CompilationPartsUtil {
runUnderStatisticsTimer(messages, 'compile-parts:upload:time') {
CompilationPartsUploader uploader = new CompilationPartsUploader(serverUrl, messages)
Set<String> alreadyUploaded = new HashSet<>()
boolean fallbackToHeads
def files = uploader.getFoundAndMissingFiles(metadataJson)
if (files != null) {
messages.info("Successfully fetched info about already uploaded files")
alreadyUploaded.addAll(files.found)
fallbackToHeads = false
}
else {
messages.warning("Failed to fetch info about already uploaded files, will fallback to HEAD requests")
fallbackToHeads = true
}
// Upload with higher threads count
executor.setMaximumPoolSize(executorThreadsCount * 2)
executor.prestartAllCoreThreads()
contexts.each { PackAndUploadContext ctx ->
if (alreadyUploaded.contains(ctx.name)) {
reusedCount.getAndIncrement()
reusedBytes.getAndAdd(new File(ctx.archive).size())
return
}
executor.submit {
def archiveFile = new File(ctx.archive)
String hash = hashes.get(ctx.name)
def path = "$uploadPrefix/${ctx.name}/${hash}.jar".toString()
if (uploader.upload(path, archiveFile)) {
if (uploader.upload(path, archiveFile, fallbackToHeads)) {
uploadedCount.getAndIncrement()
uploadedBytes.getAndAdd(archiveFile.size())
}
@@ -229,16 +256,9 @@ class CompilationPartsUtil {
executor.reportErrors(messages)
// Prepare and publish metadata file
// Save and publish metadata file
def metadataFile = new File("$zipsLocation/metadata.json")
CompilationPartsMetadata m = new CompilationPartsMetadata()
m.serverUrl = serverUrl
m.branch = branch
m.prefix = uploadPrefix
m.files = new TreeMap<String, String>(hashes)
FileUtil.writeToFile(metadataFile, new Gson().toJson(m))
FileUtil.writeToFile(metadataFile, metadataJson)
messages.artifactBuilt(metadataFile.absolutePath)
}
@@ -628,7 +648,7 @@ class CompilationPartsUtil {
futures.last().get()
}
else {
Thread.sleep(TimeUnit.SECONDS.toMillis(futures.size() < 500 ? 1 : 3))
Thread.sleep(TimeUnit.SECONDS.toMillis(1))
}
}
}