ai.timefold.solver
timefold-solver-service-jackson
diff --git a/service/definition/src/main/java/ai/timefold/solver/service/definition/impl/storage/CompressionUtils.java b/service/definition/src/main/java/ai/timefold/solver/service/definition/impl/storage/CompressionUtils.java
new file mode 100644
index 00000000000..b8f3906722c
--- /dev/null
+++ b/service/definition/src/main/java/ai/timefold/solver/service/definition/impl/storage/CompressionUtils.java
@@ -0,0 +1,104 @@
+package ai.timefold.solver.service.definition.impl.storage;
+
+import java.io.ByteArrayInputStream;
+import java.io.ByteArrayOutputStream;
+import java.io.IOException;
+import java.io.InputStream;
+import java.io.OutputStream;
+import java.io.PushbackInputStream;
+import java.util.zip.GZIPInputStream;
+import java.util.zip.GZIPOutputStream;
+
+/**
+ * Utility methods to compress and decompress content that is stored in a
+ * {@link ai.timefold.solver.service.definition.internal.storage.Storage}.
+ *
+ * Content is compressed with gzip. As storages can also contain content that was not written by the service, all
+ * decompression methods first check for the gzip magic bytes and pass the content through unchanged when it is not
+ * compressed.
+ */
+public class CompressionUtils {
+
+ private static final int GZIP_MAGIC_LENGTH = 2;
+
+ private CompressionUtils() {
+ }
+
+ public static byte[] compress(byte[] data) {
+ byte[] processed;
+ var os = new ByteArrayOutputStream();
+ try (var gzipOs = new GZIPOutputStream(os)) {
+
+ gzipOs.write(data, 0, data.length);
+ // the deflated data and the trailer only reach the underlying stream once the gzip stream is finished
+ gzipOs.finish();
+
+ processed = os.toByteArray();
+ } catch (IOException e) {
+ processed = data;
+ }
+ return processed;
+ }
+
+ public static byte[] uncompress(byte[] data) {
+ if (isCompressed(data)) {
+
+ var os = new ByteArrayOutputStream();
+
+ try (var gis = new GZIPInputStream(new ByteArrayInputStream(data))) {
+ gis.transferTo(os);
+
+ return os.toByteArray();
+ } catch (IOException e) {
+ return data;
+ }
+ } else {
+ return data;
+ }
+ }
+
+ public static boolean isCompressed(final byte[] compressed) {
+ return compressed.length >= GZIP_MAGIC_LENGTH
+ && (compressed[0] == (byte) (GZIPInputStream.GZIP_MAGIC))
+ && (compressed[1] == (byte) (GZIPInputStream.GZIP_MAGIC >> 8));
+ }
+
+ /**
+ * Wraps given stream into a decompressing one, in case its content is compressed.
+ *
+ * @param source stream to read the content from
+ * @return stream that provides the uncompressed content, closing it closes the source stream as well
+ */
+ public static InputStream decompressIfNeeded(InputStream source) throws IOException {
+ var pushback = new PushbackInputStream(source, GZIP_MAGIC_LENGTH);
+ var magic = new byte[GZIP_MAGIC_LENGTH];
+ int read = pushback.read(magic);
+ if (read > 0) {
+ pushback.unread(magic, 0, read);
+ }
+ if (read == GZIP_MAGIC_LENGTH && isCompressed(magic)) {
+ return new GZIPInputStream(pushback);
+ }
+ return pushback;
+ }
+
+ /**
+ * Transfers the content of the source stream to the target one, compressing it on the fly in case it is not
+ * compressed already.
+ *
+ * @param source stream to read the content from
+ * @param target stream to write the (compressed) content to
+ */
+ public static void transferDataCompressIfNeeded(InputStream source, OutputStream target) throws IOException {
+ byte[] twoFirst = source.readNBytes(GZIP_MAGIC_LENGTH);
+ if (isCompressed(twoFirst)) {
+ target.write(twoFirst);
+ source.transferTo(target);
+ } else {
+ var gzipOut = new GZIPOutputStream(target);
+ gzipOut.write(twoFirst);
+ source.transferTo(gzipOut);
+ gzipOut.finish();
+ }
+ }
+}
diff --git a/service/definition/src/main/java/ai/timefold/solver/service/definition/impl/storage/StorageObjectMapperProducer.java b/service/definition/src/main/java/ai/timefold/solver/service/definition/impl/storage/StorageObjectMapperProducer.java
new file mode 100644
index 00000000000..729dafbd23a
--- /dev/null
+++ b/service/definition/src/main/java/ai/timefold/solver/service/definition/impl/storage/StorageObjectMapperProducer.java
@@ -0,0 +1,37 @@
+package ai.timefold.solver.service.definition.impl.storage;
+
+import jakarta.enterprise.inject.Produces;
+import jakarta.inject.Inject;
+
+import com.fasterxml.jackson.databind.DeserializationFeature;
+import com.fasterxml.jackson.databind.MapperFeature;
+import com.fasterxml.jackson.databind.ObjectMapper;
+
+/**
+ * Produces StorageObjectMapper instances with storage-specific (de)serialization settings.
+ */
+public class StorageObjectMapperProducer {
+
+ private ObjectMapper quarkusObjectMapper;
+
+ @Inject
+ public StorageObjectMapperProducer(ObjectMapper quarkusObjectMapper) {
+ this.quarkusObjectMapper = quarkusObjectMapper;
+ }
+
+ @Produces
+ public StorageObjectMapperWrapper create() {
+ var objectMapper = quarkusObjectMapper.copy(); // Create a copy to avoid mutating the injected instance.
+ /*
+ * Storage-specific customizations are more permissive (beyond the model JSON Schema):
+ * 1. Ignore unknown properties during deserialization to allow for forward compatibility.
+ * 2. Enable case-insensitive property matching as some cloud storages are case-insensitive.
+ */
+ objectMapper.disable(DeserializationFeature.FAIL_ON_UNKNOWN_PROPERTIES);
+ objectMapper.getDeserializationConfig().with(MapperFeature.ACCEPT_CASE_INSENSITIVE_PROPERTIES, true);
+
+ objectMapper.findAndRegisterModules();
+
+ return new StorageObjectMapperWrapper(objectMapper);
+ }
+}
diff --git a/service/definition/src/main/java/ai/timefold/solver/service/definition/impl/storage/StorageObjectMapperWrapper.java b/service/definition/src/main/java/ai/timefold/solver/service/definition/impl/storage/StorageObjectMapperWrapper.java
new file mode 100644
index 00000000000..ef749c90fd5
--- /dev/null
+++ b/service/definition/src/main/java/ai/timefold/solver/service/definition/impl/storage/StorageObjectMapperWrapper.java
@@ -0,0 +1,21 @@
+package ai.timefold.solver.service.definition.impl.storage;
+
+import com.fasterxml.jackson.databind.ObjectMapper;
+
+/**
+ * Wrapper around ObjectMapper to be injected into storage services to enforce correct (de)serialization settings.
+ *
+ * Instances are normally produced by the {@link StorageObjectMapperProducer}.
+ */
+public final class StorageObjectMapperWrapper {
+
+ private final ObjectMapper objectMapper;
+
+ public StorageObjectMapperWrapper(ObjectMapper objectMapper) {
+ this.objectMapper = objectMapper;
+ }
+
+ public ObjectMapper get() {
+ return objectMapper;
+ }
+}
diff --git a/service/definition/src/main/java/ai/timefold/solver/service/definition/impl/storage/inmemory/InMemoryStorage.java b/service/definition/src/main/java/ai/timefold/solver/service/definition/impl/storage/inmemory/InMemoryStorage.java
deleted file mode 100644
index af1a9f2b705..00000000000
--- a/service/definition/src/main/java/ai/timefold/solver/service/definition/impl/storage/inmemory/InMemoryStorage.java
+++ /dev/null
@@ -1,198 +0,0 @@
-package ai.timefold.solver.service.definition.impl.storage.inmemory;
-
-import java.io.ByteArrayOutputStream;
-import java.io.IOException;
-import java.io.InputStream;
-import java.io.OutputStream;
-import java.util.List;
-import java.util.Map;
-import java.util.concurrent.ConcurrentHashMap;
-import java.util.zip.GZIPOutputStream;
-
-import ai.timefold.solver.service.definition.api.ModelOutput;
-import ai.timefold.solver.service.definition.api.domain.Metadata;
-import ai.timefold.solver.service.definition.internal.error.ErrorCodes;
-import ai.timefold.solver.service.definition.internal.error.ItemNotFoundException;
-import ai.timefold.solver.service.definition.internal.error.TimefoldRuntimeException;
-import ai.timefold.solver.service.definition.internal.storage.Storage;
-import ai.timefold.solver.service.definition.internal.storage.StorageAddress;
-import ai.timefold.solver.service.definition.internal.storage.StorageConfiguration;
-import ai.timefold.solver.service.definition.internal.storage.SubModelKind;
-
-import com.fasterxml.jackson.core.type.TypeReference;
-import com.fasterxml.jackson.databind.ObjectMapper;
-
-public abstract class InMemoryStorage implements Storage {
-
- private ObjectMapper mapper;
-
- private Map models = new ConcurrentHashMap<>();
-
- private Map resources = new ConcurrentHashMap<>();
-
- public InMemoryStorage() {
-
- }
-
- public InMemoryStorage(ObjectMapper mapper) {
- this.mapper = mapper;
- }
-
- @Override
- public void store(StorageAddress options, String id, ModelOutput_ dataset) {
- models.put(id, dataset);
- }
-
- @Override
- public void update(StorageAddress options, String id, ModelOutput_ dataset) {
- models.put(id, dataset);
- }
-
- @Override
- public void complete(StorageAddress options, String id, ModelOutput_ dataset) {
- models.put(id, dataset);
- }
-
- @Override
- public ModelOutput_ get(StorageAddress options, String id) {
- return models.get(id);
- }
-
- @Override
- public void delete(StorageAddress options, String id) {
- for (SubModelKind submodel : SubModelKind.values()) {
- Object removed = resources.remove(id + "_" + submodel.id());
-
- if (removed != null) {
- resources.put(id + "_" + submodel.id() + ".deleted", removed);
- }
- }
- ModelOutput_ dataset = models.remove(id);
- if (dataset != null) {
- models.put(id + ".deleted", dataset);
- }
- }
-
- @Override
- public void restore(StorageAddress options, String id) {
- if (!models.containsKey(id + ".deleted")) {
- throw new ItemNotFoundException(ErrorCodes.STORAGE_NO_JOB_FOUND,
- "Run with id " + id + " cannot be restored as it does not exist");
- }
- ModelOutput_ dataset = models.remove(id + ".deleted");
- models.put(id, dataset);
- for (SubModelKind submodel : SubModelKind.values()) {
- Object restored = resources.remove(id + "_" + submodel.id() + ".deleted");
-
- if (restored != null) {
- resources.put(id + "_" + submodel.id(), restored);
- }
- }
- }
-
- @Override
- public boolean exists(StorageAddress options, String id) {
- return models.containsKey(id);
- }
-
- @Override
- public List list(StorageAddress options, int pageNumber, int pageSize) {
- return resources.values().stream().filter(item -> item instanceof Metadata)
- .skip((long) pageNumber * pageSize).limit(pageSize).map(Metadata.class::cast).toList();
- }
-
- @Override
- public T getSubModel(StorageAddress options, String id, SubModelKind kind, Class clazz) {
-
- return (T) resources.get(id + "_" + kind.id());
- }
-
- @Override
- public void storeSubModel(StorageAddress options, String id, SubModelKind kind, Object subModel) {
- if (subModel == null) {
- return;
- }
- resources.put(id + "_" + kind.id(), subModel);
- }
-
- @Override
- public void storeSubModelStream(StorageAddress options, String id, SubModelKind kind, InputStream input) {
- if (input == null) {
- return;
- }
- try {
- resources.put(id + "_" + kind.id(), input.readAllBytes());
- } catch (IOException e) {
- throw new TimefoldRuntimeException(ErrorCodes.STORAGE_UNABLE_TO_WRITE,
- "Unable to store sub model (" + kind + ") to the storage for id " + id, e);
- }
- }
-
- @Override
- public void updateSubModel(StorageAddress options, String id, SubModelKind subModelKind, Object subModel) {
- if (subModel == null) {
- return;
- }
- resources.put(id + "_" + subModelKind.id(), subModel);
- }
-
- @Override
- public boolean existsSubModel(StorageAddress options, String id, SubModelKind kind) {
- return resources.containsKey(id + "_" + kind.id());
- }
-
- @Override
- public void getSubModelStream(StorageAddress options, String id, SubModelKind subModelKind, OutputStream output) {
- Object subModel = resources.get(id + "_" + subModelKind.id());
- if (subModel != null) {
- try {
- byte[] content = subModel instanceof byte[] raw ? raw : mapper.writeValueAsBytes(subModel);
- output.write(compress(content));
-
- } catch (IOException e) {
- throw new TimefoldRuntimeException(ErrorCodes.STORAGE_UNABLE_TO_READ,
- "Unable to read sub model (" + subModelKind + ") from the storage for id " + id, e);
- }
- }
- }
-
- @Override
- public void clean(StorageAddress options) {
- models.clear();
- }
-
- @Override
- public void create(String location, StorageConfiguration configuration) {
- }
-
- @Override
- public void reconfigure(String location, StorageConfiguration configuration) {
-
- }
-
- @Override
- public void destroy(String id) {
-
- }
-
- @Override
- public T getSubModel(StorageAddress options, String id, SubModelKind kind, TypeReference configurationClass) {
- return (T) resources.get(id + "_" + kind.id());
- }
-
- private byte[] compress(byte[] data) {
- byte[] processed;
- ByteArrayOutputStream os = new ByteArrayOutputStream();
- try (GZIPOutputStream gzipOs = new GZIPOutputStream(os)) {
-
- gzipOs.write(data, 0, data.length);
-
- gzipOs.close();
-
- processed = os.toByteArray();
- } catch (IOException e) {
- processed = data;
- }
- return processed;
- }
-}
diff --git a/service/definition/src/main/java/ai/timefold/solver/service/definition/internal/platform/ServiceAccountToken.java b/service/definition/src/main/java/ai/timefold/solver/service/definition/internal/platform/ServiceAccountToken.java
new file mode 100644
index 00000000000..12231b59be1
--- /dev/null
+++ b/service/definition/src/main/java/ai/timefold/solver/service/definition/internal/platform/ServiceAccountToken.java
@@ -0,0 +1,109 @@
+package ai.timefold.solver.service.definition.internal.platform;
+
+import java.io.IOException;
+import java.nio.file.Files;
+import java.nio.file.Path;
+import java.time.Duration;
+import java.util.Optional;
+
+import org.slf4j.Logger;
+import org.slf4j.LoggerFactory;
+
+/**
+ * The Kubernetes service account token of the pod, as it is mounted into it, for the clients that have to prove to
+ * the platform services who they are.
+ *
+ * The token is a JWT the cluster issues for the service account the pod runs as. It carries the namespace, the pod
+ * and the service account it was issued for, which is what lets the service on the other end tell one caller from
+ * another. The kubelet rotates the token well before it expires and rewrites the file in place, so it is read again
+ * every {@value #REFRESH_INTERVAL_SECONDS} seconds rather than kept for the lifetime of the client.
+ *
+ * Outside a cluster there is no such file, which is not an error: the token is then simply not available and the
+ * clients send their requests without one, the same way they did before the platform asked for it.
+ */
+public final class ServiceAccountToken {
+
+ private static final Logger LOGGER = LoggerFactory.getLogger(ServiceAccountToken.class);
+
+ /**
+ * File the token is read from, for the deployments that mount it somewhere else than Kubernetes does by default.
+ */
+ public static final String TOKEN_FILE_PROPERTY = "timefold.platform.service-account.token-file";
+
+ /**
+ * Where Kubernetes mounts the token of the service account the pod runs as.
+ */
+ public static final String DEFAULT_TOKEN_FILE = "/var/run/secrets/kubernetes.io/serviceaccount/token";
+
+ /**
+ * Header the token is sent in, as a bearer credential.
+ */
+ public static final String AUTHORIZATION_HEADER = "Authorization";
+
+ private static final String BEARER_PREFIX = "Bearer ";
+
+ private static final long REFRESH_INTERVAL_SECONDS = 60;
+
+ private final Path file;
+
+ private final long refreshIntervalNanos;
+
+ /**
+ * What was last read, replaced as a whole so that a reader never sees a half-written one.
+ */
+ private volatile Reading reading;
+
+ public ServiceAccountToken(String file) {
+ this(Path.of(file == null || file.isBlank() ? DEFAULT_TOKEN_FILE : file),
+ Duration.ofSeconds(REFRESH_INTERVAL_SECONDS));
+ }
+
+ public ServiceAccountToken(Path file, Duration refreshInterval) {
+ this.file = file;
+ this.refreshIntervalNanos = refreshInterval.toNanos();
+ }
+
+ /**
+ * The value of the {@value #AUTHORIZATION_HEADER} header to send, empty when there is no token to send.
+ */
+ public Optional authorization() {
+ var current = reading;
+ if (current != null && current.isFresh(System.nanoTime(), refreshIntervalNanos)) {
+ return Optional.ofNullable(current.authorization());
+ }
+ Reading read = read();
+ this.reading = read;
+ return Optional.ofNullable(read.authorization());
+ }
+
+ private Reading read() {
+ try {
+ String token = Files.readString(file).trim();
+ if (token.isEmpty()) {
+ LOGGER.warn("The service account token file {} is empty, requests are sent without a token", file);
+ return new Reading(null, System.nanoTime());
+ }
+ return new Reading(BEARER_PREFIX + token, System.nanoTime());
+ } catch (IOException | RuntimeException e) {
+ /*
+ * Only worth mentioning once in a while, as this is the normal state of anything that does not run in a
+ * cluster, and the service on the other end is the one that decides whether a caller without a token is
+ * served.
+ */
+ LOGGER.debug("No service account token could be read from {} ({}), requests are sent without one", file,
+ e.getMessage());
+ return new Reading(null, System.nanoTime());
+ }
+ }
+
+ /**
+ * @param authorization the header value that was built out of the token, null when there was no token to read
+ * @param readAtNanos when the file was last looked at
+ */
+ private record Reading(String authorization, long readAtNanos) {
+
+ boolean isFresh(long now, long refreshIntervalNanos) {
+ return now - readAtNanos < refreshIntervalNanos;
+ }
+ }
+}
diff --git a/service/definition/src/main/java/ai/timefold/solver/service/definition/internal/storage/AbstractStorageService.java b/service/definition/src/main/java/ai/timefold/solver/service/definition/internal/storage/AbstractStorageService.java
index 11a91362779..86b3f9e0afb 100644
--- a/service/definition/src/main/java/ai/timefold/solver/service/definition/internal/storage/AbstractStorageService.java
+++ b/service/definition/src/main/java/ai/timefold/solver/service/definition/internal/storage/AbstractStorageService.java
@@ -1,12 +1,17 @@
package ai.timefold.solver.service.definition.internal.storage;
+import java.io.IOException;
import java.io.InputStream;
import java.io.OutputStream;
+import java.time.Duration;
+import java.util.ArrayList;
import java.util.List;
+import java.util.Map;
import java.util.concurrent.ConcurrentHashMap;
import java.util.concurrent.TimeUnit;
import java.util.concurrent.locks.Lock;
import java.util.concurrent.locks.ReentrantLock;
+import java.util.function.Supplier;
import jakarta.enterprise.inject.Instance;
import jakarta.inject.Inject;
@@ -28,26 +33,90 @@
import ai.timefold.solver.service.definition.api.validation.LegacyValidationResult;
import ai.timefold.solver.service.definition.api.validation.ValidationBuilder;
import ai.timefold.solver.service.definition.api.validation.dto.ValidationResult;
+import ai.timefold.solver.service.definition.impl.storage.CompressionUtils;
+import ai.timefold.solver.service.definition.impl.storage.StorageObjectMapperWrapper;
import ai.timefold.solver.service.definition.impl.validation.JsonMappingError;
import ai.timefold.solver.service.definition.internal.error.ErrorCodes;
import ai.timefold.solver.service.definition.internal.error.ItemNotFoundException;
import ai.timefold.solver.service.definition.internal.error.TimefoldRuntimeException;
+import org.eclipse.microprofile.config.inject.ConfigProperty;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
+import com.fasterxml.jackson.core.JsonProcessingException;
import com.fasterxml.jackson.core.type.TypeReference;
+import com.fasterxml.jackson.databind.DatabindException;
import com.fasterxml.jackson.databind.JsonMappingException;
-
+import com.fasterxml.jackson.databind.ObjectMapper;
+
+/**
+ * Object oriented facade in front of a {@link Storage}.
+ *
+ * This service owns everything that requires the knowledge of the persisted types: it serializes objects into content
+ * and compresses it before handing it over to the storage, and decompresses and deserializes the content read back
+ * from the storage. The storage itself only ever deals with raw streams of bytes.
+ *
+ * @param type representing the data set to be solved
+ * @param type representing model specific configuration overrides
+ * @param type representing metrics of the data set to be solved
+ * @param type representing metrics of the solved data set
+ * @param type representing a solved data set
+ * @param type representing the score of a solved data set
+ * @param type representing constraint justifications of a solved data set
+ */
public abstract non-sealed class AbstractStorageService
implements StorageServiceBase {
private static final Logger LOGGER = LoggerFactory.getLogger(AbstractStorageService.class);
- private static final long LOCK_TIMEOUT_SECONDS = 60;
+ static final long LOCK_TIMEOUT_SECONDS = 60;
+
+ /**
+ * Attribute that data stores supporting attributes use to store the status of the data set. Its presence means the
+ * complete metadata can be reconstructed from the attributes, without reading the content of the sub model.
+ */
+ private static final String SOLVER_STATUS_ATTRIBUTE = "solverStatus";
+
+ /**
+ * Property that turns the caching of the model output and the metadata on and off.
+ */
+ public static final String USE_CACHE_PROPERTY = "timefold.storage.use-cache";
+
+ /**
+ * Property with the number of attempts a storage operation is given before it is failed, the first one included.
+ * One means no retrying at all.
+ *
+ * The default of five, with the default waits, keeps trying for about seven and a half seconds, which is meant to
+ * outlast the handover when the access service in front of the data store is rolled out.
+ */
+ public static final String RETRY_MAX_ATTEMPTS_PROPERTY = "timefold.storage.retry.max-attempts";
+
+ /**
+ * Property with how long to wait before the second attempt; every further wait doubles it up to the maximum.
+ */
+ public static final String RETRY_INITIAL_DELAY_PROPERTY = "timefold.storage.retry.initial-delay";
+
+ /**
+ * Property with the longest a single wait between two attempts may become.
+ *
+ * Attempts are made while holding the lock of the data set, so the whole budget has to stay well below the
+ * {@value #LOCK_TIMEOUT_SECONDS} seconds others are willing to wait for that lock.
+ */
+ public static final String RETRY_MAX_DELAY_PROPERTY = "timefold.storage.retry.max-delay";
+
+ protected Storage storage;
+
+ private ObjectMapper mapper;
+ private boolean useCache = true;
+ private final Map outputCache = new ConcurrentHashMap<>();
+ private final Map> metadataCache = new ConcurrentHashMap<>();
+
+ private int retryMaxAttempts = 5;
+ private Duration retryInitialDelay = Duration.ofMillis(500);
+ private Duration retryMaxDelay = Duration.ofSeconds(5);
- protected Storage storage;
- private final ConcurrentHashMap locks = new ConcurrentHashMap<>();
+ private final ConcurrentHashMap locks = new ConcurrentHashMap<>();
private final Lock generalLock = new ReentrantLock();
@SuppressWarnings("unused")
@@ -56,19 +125,41 @@ protected AbstractStorageService() {
}
// For tests
- protected AbstractStorageService(Storage storage) {
+ protected AbstractStorageService(Storage storage, StorageObjectMapperWrapper storageObjectMapperWrapper) {
this.storage = storage;
+ this.mapper = storageObjectMapperWrapper.get();
}
@Inject
- protected void setStorage(Instance> instance) {
+ protected void setStorage(Instance instance) {
storage = instance.get();
}
+ @Inject
+ protected void setStorageObjectMapper(Instance instance) {
+ mapper = instance.get().get();
+ }
+
+ @Inject
+ protected void setCacheEnabled(
+ @ConfigProperty(name = USE_CACHE_PROPERTY, defaultValue = "true") boolean useCache) {
+ this.useCache = useCache;
+ }
+
+ @Inject
+ protected void setRetry(
+ @ConfigProperty(name = RETRY_MAX_ATTEMPTS_PROPERTY, defaultValue = "5") int maxAttempts,
+ @ConfigProperty(name = RETRY_INITIAL_DELAY_PROPERTY, defaultValue = "PT0.5S") String initialDelay,
+ @ConfigProperty(name = RETRY_MAX_DELAY_PROPERTY, defaultValue = "PT5S") String maxDelay) {
+ this.retryMaxAttempts = Math.max(1, maxAttempts);
+ this.retryInitialDelay = Duration.parse(initialDelay);
+ this.retryMaxDelay = Duration.parse(maxDelay);
+ }
+
public void create(String id, StorageConfiguration storageConfiguration) {
acquireLock(id);
try {
- this.storage.create(id, storageConfiguration);
+ runWithRetry(id, "create the storage", () -> this.storage.create(id, storageConfiguration));
} finally {
releaseLock(id);
}
@@ -77,7 +168,7 @@ public void create(String id, StorageConfiguration storageConfiguration) {
public void reconfigure(String id, StorageConfiguration storageConfiguration) {
acquireLock(id);
try {
- this.storage.reconfigure(id, storageConfiguration);
+ runWithRetry(id, "reconfigure the storage", () -> this.storage.reconfigure(id, storageConfiguration));
} finally {
releaseLock(id);
}
@@ -86,7 +177,9 @@ public void reconfigure(String id, StorageConfiguration storageConfiguration) {
public void destroy(String id) {
acquireLock(id);
try {
- this.storage.destroy(id);
+ runWithRetry(id, "destroy the storage", () -> this.storage.destroy(id));
+ outputCache.clear();
+ metadataCache.clear();
} finally {
releaseLock(id);
}
@@ -109,17 +202,17 @@ public ModelInput_ getModelInput(String id) {
public ModelInput_ getModelInput(StorageAddress storageAddress, String id) {
acquireLock(id);
try {
- return (ModelInput_) storage.getSubModel(storageAddress, id, SubModelKind.MODEL_INPUT, getModelInputClass());
+ return (ModelInput_) readSubModel(storageAddress, id, SubModelKind.MODEL_INPUT, getModelInputClass());
} catch (TimefoldRuntimeException timefoldRuntimeException) {
if (timefoldRuntimeException.getCause() instanceof JsonMappingException mappingException) {
var validationBuilder = new ValidationBuilder().addIssue(new JsonMappingError(mappingException.getMessage()));
LegacyValidationResult legacyValidationResult = validationBuilder.buildLegacyValidationResult();
- Metadata metadata = (Metadata) storage.getSubModel(id, SubModelKind.METADATA, Metadata.class);
+ Metadata metadata = readMetadata(storageAddress, id);
if (metadata != null) {
metadata.datasetValidated(legacyValidationResult);
- storage.updateSubModel(storageAddress, id, SubModelKind.METADATA, metadata); // For backward compatibility.
- storage.storeSubModel(storageAddress, id, SubModelKind.VALIDATION_RESULT, validationBuilder.build());
+ updateSubModel(storageAddress, id, SubModelKind.METADATA, metadata); // For backward compatibility.
+ storeSubModel(storageAddress, id, SubModelKind.VALIDATION_RESULT, validationBuilder.build());
// Avoid re-throwing the exception, since it was already handled.
return null;
} else {
@@ -138,18 +231,13 @@ public ModelInput_ getModelInput(StorageAddress storageAddress, String id) {
}
public void storeModelInput(String id, ModelInput_ modelInput) {
- acquireLock(id);
- try {
- storage.storeSubModel(id, SubModelKind.MODEL_INPUT, modelInput);
- } finally {
- releaseLock(id);
- }
+ storeModelInput(null, id, modelInput);
}
public void storeModelInput(StorageAddress storageAddress, String id, ModelInput_ modelInput) {
acquireLock(id);
try {
- storage.storeSubModel(storageAddress, id, SubModelKind.MODEL_INPUT, modelInput);
+ storeSubModel(storageAddress, id, SubModelKind.MODEL_INPUT, modelInput);
} finally {
releaseLock(id);
}
@@ -162,160 +250,190 @@ public ModelInput_ getSolvedModelInput(String id) {
public ModelInput_ getSolvedModelInput(StorageAddress storageAddress, String id) {
acquireLock(id);
try {
- return (ModelInput_) storage.getSubModel(storageAddress, id, SubModelKind.MODEL_INPUT_SOLVED, getModelInputClass());
+ return (ModelInput_) readSubModel(storageAddress, id, SubModelKind.MODEL_INPUT_SOLVED, getModelInputClass());
} finally {
releaseLock(id);
}
}
public void storeSolvedModelInput(String id, ModelInput_ modelInput) {
- acquireLock(id);
- try {
- storage.storeSubModel(id, SubModelKind.MODEL_INPUT_SOLVED, modelInput);
- } finally {
- releaseLock(id);
- }
+ storeSolvedModelInput(null, id, modelInput);
}
public void storeSolvedModelInput(StorageAddress storageAddress, String id, ModelInput_ modelInput) {
acquireLock(id);
try {
- storage.storeSubModel(storageAddress, id, SubModelKind.MODEL_INPUT_SOLVED, modelInput);
+ storeSubModel(storageAddress, id, SubModelKind.MODEL_INPUT_SOLVED, modelInput);
} finally {
releaseLock(id);
}
}
public Metadata getMetadata(String id) {
- acquireLock(id);
- try {
- return (Metadata) storage.getSubModel(id, SubModelKind.METADATA, Metadata.class);
- } finally {
- releaseLock(id);
- }
+ return getMetadata(null, id);
}
public Metadata getMetadata(StorageAddress storageAddress, String id) {
acquireLock(id);
try {
- return (Metadata) storage.getSubModel(storageAddress, id, SubModelKind.METADATA, Metadata.class);
+ return readMetadata(storageAddress, id);
} finally {
releaseLock(id);
}
}
public void storeMetadata(String id, Metadata metadata) {
- acquireLock(id);
- try {
- storage.storeSubModel(id, SubModelKind.METADATA, metadata);
- } finally {
- releaseLock(id);
- }
+ storeMetadata(null, id, metadata);
}
public void storeMetadata(StorageAddress storageAddress, String id, Metadata metadata) {
acquireLock(id);
try {
- storage.storeSubModel(storageAddress, id, SubModelKind.METADATA, metadata);
+ storeSubModel(storageAddress, id, SubModelKind.METADATA, metadata);
} finally {
releaseLock(id);
}
}
public void updateMetadata(String id, Metadata metadata) {
- acquireLock(id);
- try {
- storage.updateSubModel(id, SubModelKind.METADATA, metadata);
- } finally {
- releaseLock(id);
- }
+ updateMetadata(null, id, metadata);
}
public void updateMetadata(StorageAddress storageAddress, String id, Metadata metadata) {
acquireLock(id);
try {
- storage.updateSubModel(storageAddress, id, SubModelKind.METADATA, metadata);
+ updateSubModel(storageAddress, id, SubModelKind.METADATA, metadata);
} finally {
releaseLock(id);
}
}
+ public void storeValidationResponse(String id, ValidationResult validationResult) {
+ storeValidationResponse(null, id, validationResult);
+ }
+
public void storeValidationResponse(StorageAddress storageAddress, String id,
ValidationResult validationResult) {
acquireLock(id);
try {
- storage.storeSubModel(storageAddress, id, SubModelKind.VALIDATION_RESULT, validationResult);
+ storeSubModel(storageAddress, id, SubModelKind.VALIDATION_RESULT, validationResult);
} finally {
releaseLock(id);
}
}
- public void storeValidationResponse(String id, ValidationResult validationResult) {
+ public ModelOutput_ getModelOutput(String id) {
+ return getModelOutput(null, id);
+ }
+
+ public ModelOutput_ getModelOutput(StorageAddress storageAddress, String id) {
acquireLock(id);
try {
- storage.storeSubModel(id, SubModelKind.VALIDATION_RESULT, validationResult);
+ return readModelOutput(storageAddress, id);
} finally {
releaseLock(id);
}
}
- public ModelOutput_ getModelOutput(String id) {
+ public void storeModelOutput(String id, ModelOutput_ modelOutput) {
+ storeModelOutput(null, id, modelOutput);
+ }
+
+ public void storeModelOutput(StorageAddress storageAddress, String id, ModelOutput_ modelOutput) {
acquireLock(id);
try {
- return storage.get(id);
+ StorageContent content = toContent(id, modelOutput);
+ writeWithRetry(id, "store the dataset", content, () -> storage.store(storageAddress, id, content));
+ cacheModelOutput(id, modelOutput);
} finally {
releaseLock(id);
}
}
- public void storeModelOutput(String id, ModelOutput_ modelOutput) {
+ public void updateModelOutput(String id, ModelOutput_ modelOutput) {
+ updateModelOutput(null, id, modelOutput);
+ }
+
+ public void updateModelOutput(StorageAddress storageAddress, String id, ModelOutput_ modelOutput) {
acquireLock(id);
try {
- storage.store(id, modelOutput);
+ StorageContent content = toContent(id, modelOutput);
+ writeWithRetry(id, "update the dataset", content, () -> storage.update(storageAddress, id, content));
+ cacheModelOutput(id, modelOutput);
} finally {
releaseLock(id);
}
}
- public void updateModelOutput(String id, ModelOutput_ modelOutput) {
+ public void completeModelOutput(String id, ModelOutput_ modelOutput) {
+ completeModelOutput(null, id, modelOutput);
+ }
+
+ public void completeModelOutput(StorageAddress storageAddress, String id, ModelOutput_ modelOutput) {
acquireLock(id);
try {
- storage.update(id, modelOutput);
+ StorageContent content = toContent(id, modelOutput);
+ writeWithRetry(id, "complete the dataset", content, () -> storage.complete(storageAddress, id, content));
+ cacheModelOutput(id, modelOutput);
} finally {
releaseLock(id);
}
}
+ /**
+ * Stores everything a data set starts out with as one unit; the lock is held for all of it, so that nobody
+ * observes the data set with only some of the parts written. The lock is reentrant, the methods called below take
+ * it again.
+ */
public void storeProblem(StorageAddress storageAddress, String id, ModelInput_ modelInput, Metadata metadata,
Configuration unprocessedConfiguration, Configuration configuration) {
- storeUnprocessedConfiguration(storageAddress, id, unprocessedConfiguration);
- storeModelInput(storageAddress, id, modelInput);
- storeMetadata(storageAddress, id, metadata);
- storeConfiguration(storageAddress, id, configuration);
+ acquireLock(id);
+ try {
+ storeUnprocessedConfiguration(storageAddress, id, unprocessedConfiguration);
+ storeModelInput(storageAddress, id, modelInput);
+ storeMetadata(storageAddress, id, metadata);
+ storeConfiguration(storageAddress, id, configuration);
+ } finally {
+ releaseLock(id);
+ }
}
public void storeProblem(String id, ModelInput_ modelInput, Metadata metadata,
Configuration unprocessedConfiguration, Configuration configuration) {
- storeUnprocessedConfiguration(id, unprocessedConfiguration);
- storeModelInput(id, modelInput);
- storeMetadata(id, metadata);
- storeConfiguration(id, configuration);
+ storeProblem(null, id, modelInput, metadata, unprocessedConfiguration, configuration);
}
+ /**
+ * Stores the outcome of a run as one unit, so that the model output, the metadata and the metrics of a data set
+ * never disagree with each other for a reader.
+ */
public void storeSolution(String id, ModelOutput_ modelOutput, Metadata metadata, InputMetrics_ inputMetrics,
OutputMetrics_ outputMetrics) {
- storeModelOutput(id, modelOutput);
- storeMetadata(id, metadata);
- storeInputMetrics(id, inputMetrics);
- storeOutputMetrics(id, outputMetrics);
+ acquireLock(id);
+ try {
+ storeModelOutput(id, modelOutput);
+ storeMetadata(id, metadata);
+ storeInputMetrics(id, inputMetrics);
+ storeOutputMetrics(id, outputMetrics);
+ } finally {
+ releaseLock(id);
+ }
}
+ /**
+ * Updates the outcome of a run as one unit; see {@link #storeSolution}.
+ */
public void updateSolution(String id, ModelOutput_ modelOutput, Metadata metadata, InputMetrics_ inputMetrics,
OutputMetrics_ outputMetrics) {
- updateModelOutput(id, modelOutput);
- updateMetadata(id, metadata);
- updateInputMetrics(id, inputMetrics);
- updateOutputMetrics(id, outputMetrics);
+ acquireLock(id);
+ try {
+ updateModelOutput(id, modelOutput);
+ updateMetadata(id, metadata);
+ updateInputMetrics(id, inputMetrics);
+ updateOutputMetrics(id, outputMetrics);
+ } finally {
+ releaseLock(id);
+ }
}
public Configuration getConfiguration(String id) {
@@ -325,7 +443,7 @@ public Configuration getConfiguration(String id) {
public Configuration getConfiguration(StorageAddress storageAddress, String id) {
acquireLock(id);
try {
- return storage.getSubModel(storageAddress, id, SubModelKind.CONFIG, getConfigurationClass());
+ return readSubModel(storageAddress, id, SubModelKind.CONFIG, getConfigurationClass());
} finally {
releaseLock(id);
}
@@ -338,71 +456,64 @@ public Configuration getUnprocessedConfiguration(String i
public Configuration getUnprocessedConfiguration(StorageAddress storageAddress, String id) {
acquireLock(id);
try {
- return storage.getSubModel(storageAddress, id, SubModelKind.UNPROCESSED_CONFIG, getConfigurationClass());
+ return readSubModel(storageAddress, id, SubModelKind.UNPROCESSED_CONFIG, getConfigurationClass());
} finally {
releaseLock(id);
}
}
public void storeConfiguration(String id, Configuration configuration) {
- acquireLock(id);
- try {
- storage.storeSubModel(id, SubModelKind.CONFIG, configuration);
- } finally {
- releaseLock(id);
- }
+ storeConfiguration(null, id, configuration);
}
public void storeConfiguration(StorageAddress storageAddress, String id,
Configuration configuration) {
acquireLock(id);
try {
- storage.storeSubModel(storageAddress, id, SubModelKind.CONFIG, configuration);
+ storeSubModel(storageAddress, id, SubModelKind.CONFIG, configuration);
} finally {
releaseLock(id);
}
}
public void storeUnprocessedConfiguration(String id, Configuration configuration) {
- acquireLock(id);
- try {
- storage.storeSubModel(id, SubModelKind.UNPROCESSED_CONFIG, configuration);
- } finally {
- releaseLock(id);
- }
+ storeUnprocessedConfiguration(null, id, configuration);
}
public void storeUnprocessedConfiguration(StorageAddress storageAddress, String id,
Configuration configuration) {
acquireLock(id);
try {
- storage.storeSubModel(storageAddress, id, SubModelKind.UNPROCESSED_CONFIG, configuration);
+ storeSubModel(storageAddress, id, SubModelKind.UNPROCESSED_CONFIG, configuration);
} finally {
releaseLock(id);
}
}
+ /**
+ * Reads the parts the solver starts from as one unit, so that they all come from the same state of the data set.
+ */
public SolverInput getSolverInput(String id) {
- ModelInput_ modelInput = getModelInput(id);
- ModelOutput_ modelOutput = getModelOutput(id);
- Configuration configuration = getConfiguration(id);
- return new SolverInput<>(modelInput, configuration, modelOutput);
- }
-
- public void storePatchRequest(String id, ModelInputPatchRequest patchRequest) {
acquireLock(id);
try {
- storage.storeSubModel(id, SubModelKind.PATCH_REQUEST, patchRequest);
+ ModelInput_ modelInput = getModelInput(id);
+ ModelOutput_ modelOutput = getModelOutput(id);
+ Configuration configuration = getConfiguration(id);
+ return new SolverInput<>(modelInput, configuration, modelOutput);
} finally {
releaseLock(id);
}
}
+ public void storePatchRequest(String id, ModelInputPatchRequest patchRequest) {
+ storePatchRequest(null, id, patchRequest);
+ }
+
public void storePatchRequest(StorageAddress storageAddress, String id,
ModelInputPatchRequest patchRequest) {
acquireLock(id);
try {
- storage.storeSubModel(storageAddress, id, SubModelKind.PATCH_REQUEST, patchRequest);
+ storeSubModel(storageAddress, id, SubModelKind.PATCH_REQUEST, patchRequest);
} finally {
releaseLock(id);
}
@@ -416,15 +527,15 @@ public ModelResponse getMod
String id) {
acquireLock(id);
try {
- Metadata metadata = storage.getSubModel(storageAddress, id, SubModelKind.METADATA, Metadata.class);
+ Metadata metadata = readMetadata(storageAddress, id);
if (metadata == null) {
throw new ItemNotFoundException(ErrorCodes.STORAGE_NO_JOB_FOUND, "Unable to find data set for id " + id);
}
try {
- ModelOutput_ modelOutput = storage.get(storageAddress, id);
+ ModelOutput_ modelOutput = readModelOutput(storageAddress, id);
OutputMetrics_ outputMetrics =
- (OutputMetrics_) storage.getSubModel(storageAddress, id, SubModelKind.KPIS, getOutputMetricsClass());
- InputMetrics_ inputMetrics = (InputMetrics_) storage.getSubModel(storageAddress, id, SubModelKind.INPUT_METRICS,
+ (OutputMetrics_) readSubModel(storageAddress, id, SubModelKind.KPIS, getOutputMetricsClass());
+ InputMetrics_ inputMetrics = (InputMetrics_) readSubModel(storageAddress, id, SubModelKind.INPUT_METRICS,
getInputMetricsClass());
return new ModelResponse<>(metadata, modelOutput, inputMetrics, outputMetrics);
} catch (ItemNotFoundException e) {
@@ -440,20 +551,29 @@ public ModelRequest getModelRequest(String i
return getModelRequest(null, id);
}
+ /**
+ * Reads the request the data set was created from as one unit, so that the input and the configuration come from
+ * the same state of the data set.
+ */
public ModelRequest getModelRequest(StorageAddress storageAddress, String id) {
- ModelInput_ modelInput = getModelInput(storageAddress, id);
- Configuration unprocessedConfiguration = getUnprocessedConfiguration(storageAddress, id);
- if (unprocessedConfiguration != null) {
- return new ModelRequest<>(unprocessedConfiguration, modelInput);
+ acquireLock(id);
+ try {
+ ModelInput_ modelInput = getModelInput(storageAddress, id);
+ Configuration unprocessedConfiguration = getUnprocessedConfiguration(storageAddress, id);
+ if (unprocessedConfiguration != null) {
+ return new ModelRequest<>(unprocessedConfiguration, modelInput);
+ }
+ Configuration configOverrides = getConfiguration(storageAddress, id);
+ return new ModelRequest<>(configOverrides, modelInput);
+ } finally {
+ releaseLock(id);
}
- Configuration configOverrides = getConfiguration(storageAddress, id);
- return new ModelRequest<>(configOverrides, modelInput);
}
public void storeInputMetrics(String id, InputMetrics_ inputMetrics) {
acquireLock(id);
try {
- storage.storeSubModel(id, SubModelKind.INPUT_METRICS, inputMetrics);
+ storeSubModel(null, id, SubModelKind.INPUT_METRICS, inputMetrics);
} finally {
releaseLock(id);
}
@@ -462,7 +582,7 @@ public void storeInputMetrics(String id, InputMetrics_ inputMetrics) {
public void storeOutputMetrics(String id, OutputMetrics_ outputMetrics) {
acquireLock(id);
try {
- storage.storeSubModel(id, SubModelKind.KPIS, outputMetrics);
+ storeSubModel(null, id, SubModelKind.KPIS, outputMetrics);
} finally {
releaseLock(id);
}
@@ -471,7 +591,7 @@ public void storeOutputMetrics(String id, OutputMetrics_ outputMetrics) {
public void updateInputMetrics(String id, InputMetrics_ inputMetrics) {
acquireLock(id);
try {
- storage.updateSubModel(id, SubModelKind.INPUT_METRICS, inputMetrics);
+ updateSubModel(null, id, SubModelKind.INPUT_METRICS, inputMetrics);
} finally {
releaseLock(id);
}
@@ -480,43 +600,33 @@ public void updateInputMetrics(String id, InputMetrics_ inputMetrics) {
public void updateOutputMetrics(String id, OutputMetrics_ outputMetrics) {
acquireLock(id);
try {
- storage.updateSubModel(id, SubModelKind.KPIS, outputMetrics);
+ updateSubModel(null, id, SubModelKind.KPIS, outputMetrics);
} finally {
releaseLock(id);
}
}
public LogInfo getLogs(String id) {
- acquireLock(id);
- try {
- return storage.getSubModel(id, SubModelKind.LOGS, LogInfo.class);
- } finally {
- releaseLock(id);
- }
+ return getLogs(null, id);
}
public LogInfo getLogs(StorageAddress storageAddress, String id) {
acquireLock(id);
try {
- return storage.getSubModel(storageAddress, id, SubModelKind.LOGS, LogInfo.class);
+ return readSubModel(storageAddress, id, SubModelKind.LOGS, LogInfo.class);
} finally {
releaseLock(id);
}
}
public void storeLogs(String id, LogInfo info) {
- acquireLock(id);
- try {
- storage.storeSubModel(id, SubModelKind.LOGS, info);
- } finally {
- releaseLock(id);
- }
+ storeLogs(null, id, info);
}
public void storeLogs(StorageAddress storageAddress, String id, LogInfo info) {
acquireLock(id);
try {
- storage.storeSubModel(storageAddress, id, SubModelKind.LOGS, info);
+ storeSubModel(storageAddress, id, SubModelKind.LOGS, info);
} finally {
releaseLock(id);
}
@@ -529,7 +639,10 @@ public void storeExecutionArtifacts(String id, InputStream input) {
public void storeExecutionArtifacts(StorageAddress storageAddress, String id, InputStream input) {
acquireLock(id);
try {
- storage.storeSubModelStream(storageAddress, id, SubModelKind.EXECUTION_ARTIFACTS, input);
+ SubModelKind kind = SubModelKind.EXECUTION_ARTIFACTS;
+ var content = new StorageContent(() -> input, -1, Map.of(), false);
+ writeWithRetry(id, "store the submodel (" + kind + ")", content,
+ () -> storage.storeSubModel(storageAddress, id, kind, content));
} finally {
releaseLock(id);
}
@@ -542,32 +655,61 @@ public void deleteAll(String id) {
public void deleteAll(StorageAddress storageAddress, String id) {
acquireLock(id);
try {
- storage.delete(storageAddress, id);
+ runWithRetry(id, "delete the dataset", () -> storage.delete(storageAddress, id));
+ outputCache.remove(id);
+ metadataCache.remove(id);
} finally {
releaseLock(id);
}
}
- protected abstract Class> getModelInputClass();
+ public void restoreAll(String id) {
+ restoreAll(null, id);
+ }
- protected abstract Class> getInputMetricsClass();
+ public void restoreAll(StorageAddress storageAddress, String id) {
+ acquireLock(id);
+ try {
+ runWithRetry(id, "restore the dataset", () -> storage.restore(storageAddress, id));
+ } finally {
+ releaseLock(id);
+ }
+ }
- protected abstract Class> getOutputMetricsClass();
+ public boolean exists(String id) {
+ return exists(null, id);
+ }
- protected abstract TypeReference> getConfigurationClass();
+ public boolean exists(StorageAddress storageAddress, String id) {
+ return withRetry(id, "check whether the dataset exists", () -> storage.exists(storageAddress, id));
+ }
+
+ public boolean existsSubModel(StorageAddress options, String id, SubModelKind subModelKind) {
+ return withRetry(id, "check whether the submodel (" + subModelKind + ") exists",
+ () -> storage.existsSubModel(options, id, subModelKind));
+ }
public List> listRuns(int pageNumber, int pageSize) {
- return storage.list(pageNumber, pageSize);
+ return listRuns(null, pageNumber, pageSize);
}
public List> listRuns(StorageAddress storageAddress, int pageNumber, int pageSize) {
- return storage.list(storageAddress, pageNumber, pageSize);
+ List items =
+ withRetry(null, "list the datasets", () -> storage.list(storageAddress, pageNumber, pageSize));
+ List> runs = new ArrayList<>(items.size());
+ for (StorageItem item : items) {
+ Metadata metadata = toMetadata(storageAddress, item);
+ if (metadata != null) {
+ runs.add(metadata);
+ }
+ }
+ return runs;
}
public T getWaypoints(String id, TypeReference clazz) {
acquireLock(id);
try {
- return storage.getSubModel(id, SubModelKind.WAYPOINTS, clazz);
+ return readSubModel(null, id, SubModelKind.WAYPOINTS, clazz);
} finally {
releaseLock(id);
}
@@ -576,7 +718,7 @@ public T getWaypoints(String id, TypeReference clazz) {
public void storeWaypoints(String id, Object waypoints) {
acquireLock(id);
try {
- storage.storeSubModel(id, SubModelKind.WAYPOINTS, waypoints);
+ storeSubModel(null, id, SubModelKind.WAYPOINTS, waypoints);
} finally {
releaseLock(id);
}
@@ -585,40 +727,397 @@ public void storeWaypoints(String id, Object waypoints) {
public T getSubModel(StorageAddress options, String id, SubModelKind config, Class clazz) {
acquireLock(id);
try {
- return storage.getSubModel(options, id, config, clazz);
+ return readSubModel(options, id, config, clazz);
+ } finally {
+ releaseLock(id);
+ }
+ }
+
+ public T getSubModel(StorageAddress options, String id, SubModelKind config, TypeReference clazz) {
+ acquireLock(id);
+ try {
+ return readSubModel(options, id, config, clazz);
+ } finally {
+ releaseLock(id);
+ }
+ }
+
+ /**
+ * Writes the content of the sub model, as it is stored, into given output stream. The content is always written
+ * compressed, regardless of how it is stored in the underlying data store.
+ * This is the one read that is not attempted again after a failure; bytes may already have reached the
+ * stream of the caller by then, and starting over would append the content twice.
+ *
+ * @param options storage address to apply during the operation, can be null to use the default location
+ * @param id unique identifier of the data set
+ * @param subModelKind kind of the sub model e.g. waypoints
+ * @param out output stream the content should be written to
+ * @throws ItemNotFoundException in case there is no such sub model
+ */
+ public void getSubModelStream(StorageAddress options, String id, SubModelKind subModelKind, OutputStream out) {
+ acquireLock(id);
+ try (InputStream content = storage.getSubModel(options, id, subModelKind)) {
+ if (content == null) {
+ throw new ItemNotFoundException(ErrorCodes.STORAGE_NO_JOB_FOUND, "Unable to find dataset for id " + id);
+ }
+ CompressionUtils.transferDataCompressIfNeeded(content, out);
+ } catch (IOException e) {
+ throw new TimefoldRuntimeException(ErrorCodes.STORAGE_UNABLE_TO_READ,
+ "Unable to read submodel (" + subModelKind + ") from the storage for id " + id, e);
} finally {
releaseLock(id);
}
}
+ protected abstract Class> getModelInputClass();
+
+ protected abstract Class> getModelOutputClass();
+
+ protected abstract Class> getInputMetricsClass();
+
+ protected abstract Class> getOutputMetricsClass();
+
+ protected abstract TypeReference> getConfigurationClass();
+
+ /*
+ * (De)serialization of the content exchanged with the storage. None of the methods below acquires the lock, they
+ * are expected to be called from methods that already hold it.
+ */
+
+ /**
+ * Reads the data set from the storage and deserializes it into the model output type.
+ *
+ * @param address storage address to apply during the operation, can be null to use the default location
+ * @param id unique identifier of the data set
+ * @return the model output, never null
+ * @throws ItemNotFoundException in case given data set does not exist
+ */
+ protected ModelOutput_ readModelOutput(StorageAddress address, String id) {
+ if (useCache) {
+ ModelOutput_ cached = outputCache.get(id);
+ if (cached != null) {
+ return cached;
+ }
+ }
+ ModelOutput_ modelOutput = (ModelOutput_) read(id, "dataset",
+ () -> storage.get(address, id), (m, content) -> m.readValue(content, getModelOutputClass()));
+ cacheModelOutput(id, modelOutput);
+ return modelOutput;
+ }
+
+ /**
+ * Reads the metadata sub model of given data set, consulting the cache first.
+ *
+ * @param address storage address to apply during the operation, can be null to use the default location
+ * @param id unique identifier of the data set
+ * @return the metadata or null if there is none
+ */
+ protected Metadata readMetadata(StorageAddress address, String id) {
+ if (useCache) {
+ Metadata cached = metadataCache.get(id);
+ if (cached != null) {
+ return cached;
+ }
+ }
+ return readSubModel(address, id, SubModelKind.METADATA, Metadata.class);
+ }
+
+ /**
+ * Reads given sub model from the storage and deserializes it into given type.
+ *
+ * @param address storage address to apply during the operation, can be null to use the default location
+ * @param id unique identifier of the data set
+ * @param kind kind of the sub model e.g. waypoints
+ * @param clazz class the sub model should be deserialized into
+ * @return the sub model or null if there is none
+ */
+ protected T readSubModel(StorageAddress address, String id, SubModelKind kind, Class clazz) {
+ return read(id, "submodel (" + kind + ")", () -> storage.getSubModel(address, id, kind),
+ (m, content) -> m.readValue(content, clazz));
+ }
+
+ /**
+ * Reads given sub model from the storage and deserializes it into given generic type.
+ *
+ * @param address storage address to apply during the operation, can be null to use the default location
+ * @param id unique identifier of the data set
+ * @param kind kind of the sub model e.g. waypoints
+ * @param typeReference type the sub model should be deserialized into
+ * @return the sub model or null if there is none
+ */
+ protected T readSubModel(StorageAddress address, String id, SubModelKind kind, TypeReference typeReference) {
+ return read(id, "submodel (" + kind + ")", () -> storage.getSubModel(address, id, kind),
+ (m, content) -> m.readValue(content, typeReference));
+ }
+
+ /**
+ * Serializes given sub model and stores it in the storage.
+ *
+ * @param address storage address to apply during the operation, can be null to use the default location
+ * @param id unique identifier of the data set
+ * @param kind kind of the sub model e.g. waypoints
+ * @param subModel the sub model to be stored
+ */
+ protected void storeSubModel(StorageAddress address, String id, SubModelKind kind, Object subModel) {
+ StorageContent content = toContent(id, subModel);
+ writeWithRetry(id, "store the submodel (" + kind + ")", content,
+ () -> storage.storeSubModel(address, id, kind, content));
+ cacheMetadata(id, kind, subModel);
+ }
+
+ /**
+ * Serializes given sub model and updates it in the storage.
+ *
+ * @param address storage address to apply during the operation, can be null to use the default location
+ * @param id unique identifier of the data set
+ * @param kind kind of the sub model e.g. waypoints
+ * @param subModel the sub model to be stored
+ */
+ protected void updateSubModel(StorageAddress address, String id, SubModelKind kind, Object subModel) {
+ StorageContent content = toContent(id, subModel);
+ writeWithRetry(id, "update the submodel (" + kind + ")", content,
+ () -> storage.updateSubModel(address, id, kind, content));
+ cacheMetadata(id, kind, subModel);
+ }
+
+ /**
+ * Serializes and compresses given value into content that can be handed over to the storage.
+ *
+ * @param id unique identifier of the data set the value belongs to
+ * @param value the value to be serialized
+ * @return content to be stored
+ */
+ protected StorageContent toContent(String id, Object value) {
+ return StorageContent.of(CompressionUtils.compress(writeAsBytes(id, value)), attributes(value));
+ }
+
+ private byte[] writeAsBytes(String id, Object value) {
+ try {
+ return mapper().writeValueAsBytes(value);
+ } catch (JsonProcessingException e) {
+ throw new TimefoldRuntimeException(ErrorCodes.STORAGE_UNABLE_TO_WRITE,
+ "Unable to write dataset to the storage for id " + id, e, false);
+ }
+ }
+
+ /**
+ * Every read opens a new stream, so a read that failed recoverably is simply attempted again.
+ */
+ private T read(String id, String description, ContentSupplier supplier, ContentReader reader) {
+ return withRetry(id, "read the " + description, () -> readOnce(id, description, supplier, reader));
+ }
+
+ private T readOnce(String id, String description, ContentSupplier supplier, ContentReader reader) {
+ try (InputStream stored = supplier.get()) {
+ if (stored == null) {
+ return null;
+ }
+ try (InputStream content = CompressionUtils.decompressIfNeeded(stored)) {
+ return reader.read(mapper(), content);
+ }
+ } catch (DatabindException e) {
+ throw new TimefoldRuntimeException(ErrorCodes.STORAGE_UNABLE_TO_READ,
+ "Unable to read " + description + " from the storage for id " + id, e, false);
+ } catch (IOException e) {
+ throw new TimefoldRuntimeException(ErrorCodes.STORAGE_UNABLE_TO_READ,
+ "Unable to read " + description + " from the storage for id " + id, e);
+ }
+ }
+
+ /**
+ * Reconstructs the metadata of given item, either from the attributes of the underlying data store or, when the
+ * data store does not provide them, from the content of the metadata sub model.
+ */
+ private Metadata toMetadata(StorageAddress address, StorageItem item) {
+ if (hasSolverStatus(item.attributes())) {
+ return mapper().convertValue(item.attributes(), Metadata.class);
+ }
+ /*
+ * The attributes were not populated when the data set was stored, fall back to reading the content. Listing
+ * itself is not done under a lock, but reading and rewriting the metadata of a single data set is, so that it
+ * does not interleave with whoever is writing that same data set.
+ */
+ acquireLock(item.id());
+ try {
+ Metadata metadata = readMetadata(address, item.id());
+ if (metadata == null) {
+ LOGGER.debug("Unable to load run with id {}", item.id());
+ return null;
+ }
+ if (storage.supportsAttributes()) {
+ // Store it again, so the attributes get populated and the content does not have to be read next time.
+ updateSubModel(address, item.id(), SubModelKind.METADATA, metadata);
+ }
+ return metadata;
+ } finally {
+ releaseLock(item.id());
+ }
+ }
+
+ private static boolean hasSolverStatus(Map attributes) {
+ // Some data stores lowercase the attribute names.
+ return attributes.keySet().stream().anyMatch(SOLVER_STATUS_ATTRIBUTE::equalsIgnoreCase);
+ }
+
+ private static Map attributes(Object value) {
+ return value instanceof Metadata> metadata ? metadata.asMap() : Map.of();
+ }
+
+ private void cacheModelOutput(String id, ModelOutput_ modelOutput) {
+ if (useCache && modelOutput != null) {
+ outputCache.put(id, modelOutput);
+ }
+ }
+
+ private void cacheMetadata(String id, SubModelKind kind, Object subModel) {
+ if (useCache && kind == SubModelKind.METADATA && subModel instanceof Metadata) {
+ metadataCache.put(id, (Metadata) subModel);
+ }
+ }
+
+ protected ObjectMapper mapper() {
+ if (mapper == null) {
+ throw new IllegalStateException("The storage object mapper has not been set");
+ }
+ return mapper;
+ }
+
+ /*
+ * Retrying of the operations against the storage. A data store, or the access service in front of it, can be
+ * briefly unavailable; since every write is a full overwrite keyed by the id of the data set, sending it again is
+ * harmless even when the failed attempt did reach the data store after all.
+ */
+
+ /**
+ * Runs a read against the storage, retrying it while it keeps failing recoverably. Reads can always be attempted
+ * again, they produce a new stream every time.
+ */
+ protected T withRetry(String id, String description, Supplier operation) {
+ return attempt(id, description, true, operation);
+ }
+
+ /**
+ * Runs an operation against the storage that has no content of its own, such as a delete or an administrative
+ * one, retrying it while it keeps failing recoverably.
+ */
+ protected void runWithRetry(String id, String description, Runnable operation) {
+ attempt(id, description, true, () -> {
+ operation.run();
+ return null;
+ });
+ }
+
+ /**
+ * Writes content to the storage, retrying it while it keeps failing recoverably. Only content that can be read
+ * again is retried; a stream somebody else opened has already been consumed by the failed attempt.
+ */
+ protected void writeWithRetry(String id, String description, StorageContent content, Runnable operation) {
+ attempt(id, description, content.repeatable(), () -> {
+ operation.run();
+ return null;
+ });
+ }
+
+ private T attempt(String id, String description, boolean repeatable, Supplier operation) {
+ var attempt = 1;
+ Duration delay = retryInitialDelay;
+ while (true) {
+ try {
+ return operation.get();
+ } catch (RuntimeException e) {
+ if (attempt >= retryMaxAttempts || !repeatable || !isRecoverable(e)) {
+ throw e;
+ }
+ LOGGER.warn("Attempt {} of {} to {} for id {} failed, retrying in {}: {}",
+ attempt, retryMaxAttempts, description, id, delay, e.getMessage());
+ await(delay, description, id);
+ delay = nextDelay(delay);
+ attempt++;
+ }
+ }
+ }
+
+ /**
+ * Only a failure the storage itself reported as recoverable is attempted again; anything else, a data set that
+ * does not exist or content that cannot be deserialized included, would fail exactly the same way next time.
+ */
+ private static boolean isRecoverable(RuntimeException e) {
+ return e instanceof TimefoldRuntimeException timefoldException && timefoldException.isRecoverable();
+ }
+
+ private Duration nextDelay(Duration delay) {
+ Duration doubled = delay.multipliedBy(2);
+ return doubled.compareTo(retryMaxDelay) > 0 ? retryMaxDelay : doubled;
+ }
+
+ private static void await(Duration delay, String description, String id) {
+ try {
+ Thread.sleep(delay.toMillis());
+ } catch (InterruptedException e) {
+ Thread.currentThread().interrupt();
+ throw new TimefoldRuntimeException(ErrorCodes.STORAGE_UNKNOWN,
+ "Interrupted while waiting to " + description + " for id " + id + " again", e, true);
+ }
+ }
+
protected void acquireLock(String id) {
- ReentrantLock idLock;
+ LockReference reference;
tryAcquireLock(generalLock);
try {
- idLock = locks.computeIfAbsent(id, k -> new ReentrantLock());
+ reference = locks.computeIfAbsent(id, k -> new LockReference());
+ reference.users++;
} finally {
generalLock.unlock();
}
- tryAcquireLock(idLock);
+ try {
+ tryAcquireLock(reference.lock);
+ } catch (RuntimeException e) {
+ // The lock was counted as taken before it was awaited, so the count has to be given back.
+ dropReference(id);
+ throw e;
+ }
}
protected void releaseLock(String id) {
tryAcquireLock(generalLock);
try {
- ReentrantLock idLock = locks.get(id);
- if (idLock != null && idLock.isHeldByCurrentThread()) {
- idLock.unlock();
- if (idLock.getQueueLength() == 0) {
- locks.remove(id, idLock);
- }
+ LockReference reference = locks.get(id);
+ if (reference != null && reference.lock.isHeldByCurrentThread()) {
+ reference.lock.unlock();
+ dropReference(id, reference);
}
} finally {
generalLock.unlock();
}
}
+ private void dropReference(String id) {
+ tryAcquireLock(generalLock);
+ try {
+ dropReference(id, locks.get(id));
+ } finally {
+ generalLock.unlock();
+ }
+ }
+
+ /**
+ * Forgets the lock of given data set once nobody is using it any more. Must be called while holding the general
+ * lock, which is also what {@link #acquireLock(String)} takes to hand the lock out, so that a lock is never
+ * dropped from under a thread that is about to take it or that still holds it. The lock is reentrant, an outer
+ * operation of the same thread counts as a user of its own.
+ */
+ private void dropReference(String id, LockReference reference) {
+ if (reference == null) {
+ return;
+ }
+ reference.users--;
+ if (reference.users <= 0) {
+ locks.remove(id, reference);
+ }
+ }
+
protected void tryAcquireLock(Lock lock) {
try {
if (!lock.tryLock(LOCK_TIMEOUT_SECONDS, TimeUnit.SECONDS)) {
@@ -626,36 +1125,33 @@ protected void tryAcquireLock(Lock lock) {
true);
}
} catch (InterruptedException e) {
- lock.unlock();
+ // The lock was never taken, there is nothing to unlock here.
Thread.currentThread().interrupt();
- throw new TimefoldRuntimeException(ErrorCodes.STORAGE_UNKNOWN, "Timeout acquiring lock to access to storage", e,
- true);
+ throw new TimefoldRuntimeException(ErrorCodes.STORAGE_UNKNOWN, "Interrupted acquiring lock to access to storage",
+ e, true);
}
}
- public void getSubModelStream(StorageAddress options, String id, SubModelKind subModelKind, OutputStream out) {
- acquireLock(id);
- try {
- storage.getSubModelStream(options, id, subModelKind, out);
- } finally {
- releaseLock(id);
- }
- }
+ /**
+ * The lock of one data set together with the number of operations that are holding or awaiting it. The counter is
+ * only ever touched while holding the general lock.
+ */
+ private static final class LockReference {
- public boolean existsSubModel(StorageAddress options, String id, SubModelKind subModelKind) {
- return storage.existsSubModel(options, id, subModelKind);
+ private final ReentrantLock lock = new ReentrantLock();
+
+ private int users;
}
- public void restoreAll(String id) {
- deleteAll(null, id);
+ @FunctionalInterface
+ private interface ContentSupplier {
+
+ InputStream get();
}
- public void restoreAll(StorageAddress storageAddress, String id) {
- acquireLock(id);
- try {
- storage.restore(storageAddress, id);
- } finally {
- releaseLock(id);
- }
+ @FunctionalInterface
+ private interface ContentReader {
+
+ T read(ObjectMapper mapper, InputStream content) throws IOException;
}
}
diff --git a/service/definition/src/main/java/ai/timefold/solver/service/definition/internal/storage/Storage.java b/service/definition/src/main/java/ai/timefold/solver/service/definition/internal/storage/Storage.java
index ad7d73c7f60..e9928de064b 100644
--- a/service/definition/src/main/java/ai/timefold/solver/service/definition/internal/storage/Storage.java
+++ b/service/definition/src/main/java/ai/timefold/solver/service/definition/internal/storage/Storage.java
@@ -1,343 +1,166 @@
package ai.timefold.solver.service.definition.internal.storage;
import java.io.InputStream;
-import java.io.OutputStream;
import java.util.List;
-import ai.timefold.solver.service.definition.api.domain.Metadata;
import ai.timefold.solver.service.definition.internal.error.ItemNotFoundException;
-import com.fasterxml.jackson.core.type.TypeReference;
-
/**
- * Storage responsible for persisting ModelOutput_ and its sub resources into a data store
- *
- * @param representing a solved data set
+ * Storage responsible for persisting data sets and their sub models into a data store.
+ *
+ * The storage deals with raw content only, it is completely unaware of the types it persists.
+ * Turning objects into content and back, including compression and decompression, is the responsibility of
+ * {@link AbstractStorageService}. Content is always exchanged as streams so that potentially large data sets
+ * do not have to be kept in memory in their entirety.
*/
-public interface Storage {
-
- public static final String DATASETS_PREFIX = "datasets";
+public interface Storage {
- public static final String RUNS_PREFIX = "run";
-
- /**
- * Stores given data set into default location in the underlying data store
- *
- * @param id unique identifier of the data set
- * @param dataset data set to be stored
- */
- default void store(String id, ModelOutput_ dataset) {
- store(null, id, dataset);
- }
+ String DATASETS_PREFIX = "datasets";
- /**
- * Stores given data set into location defined by StorageOptions
- *
- * @param options storage option to apply during the operation
- * @param id unique identifier of the data set
- * @param dataset data set to be stored
- */
- void store(StorageAddress options, String id, ModelOutput_ dataset);
+ String RUNS_PREFIX = "run";
/**
- * Updates existing data set in the storage under default location in the underlying data store
+ * Stores given data set into location defined by StorageAddress
*
+ * @param address storage address to apply during the operation, can be null to use the default location
* @param id unique identifier of the data set
- * @param dataset data set to be stored
+ * @param content content of the data set to be stored
*/
- default void update(String id, ModelOutput_ dataset) {
- update(null, id, dataset);
- }
+ void store(StorageAddress address, String id, StorageContent content);
/**
- * Updates existing data set in the storage into location defined by StorageOptions
+ * Updates existing data set in the storage in location defined by StorageAddress
*
- * @param options storage option to apply during the operation
+ * @param address storage address to apply during the operation, can be null to use the default location
* @param id unique identifier of the data set
- * @param dataset data set to be stored
+ * @param content content of the data set to be stored
*/
- void update(StorageAddress options, String id, ModelOutput_ dataset);
+ void update(StorageAddress address, String id, StorageContent content);
/**
- * Final update of the data set upon completion of solving into default location in the underlying data store
+ * Final update of the data set upon completion of solving into location defined by StorageAddress
*
+ * @param address storage address to apply during the operation, can be null to use the default location
* @param id unique identifier of the data set
- * @param dataset data set to be stored
+ * @param content content of the data set to be stored
*/
- default void complete(String id, ModelOutput_ dataset) {
- complete(null, id, dataset);
- }
+ void complete(StorageAddress address, String id, StorageContent content);
/**
- * Final update of the data set upon completion of solving into location defined by StorageOptions
+ * Retrieves content of the data set by its unique identifier from location defined by StorageAddress.
+ *
+ * It is the responsibility of the caller to close the returned stream.
*
- * @param options storage option to apply during the operation
+ * @param address storage address to apply during the operation, can be null to use the default location
* @param id unique identifier of the data set
- * @param dataset data set to be stored
- */
- void complete(StorageAddress options, String id, ModelOutput_ dataset);
-
- /**
- * Retrieves data set by its unique identifier from the default location in the underlying data store
- *
- * @param id unique identifier of the data set
- * @return loaded data set if found
+ * @return non null stream with the content of the data set
* @throws ItemNotFoundException in case given data set does not exist
*/
- default ModelOutput_ get(String id) {
- return get(null, id);
- }
+ InputStream get(StorageAddress address, String id);
/**
- * Retrieves data set by its unique identifier from location defined by StorageOptions
+ * Deletes data set with given identifier, including all its sub models, from location defined by
+ * StorageAddress
*
- * @param options storage option to apply during the operation
+ * @param address storage address to apply during the operation, can be null to use the default location
* @param id unique identifier of the data set
- * @return loaded data set if found
- * @throws ItemNotFoundException in case given data set does not exist
*/
- ModelOutput_ get(StorageAddress options, String id);
-
- /**
- * Deletes data set with given identifier from the default location in the underlying data store
- *
- * @param id unique identifier of the data set
- */
- default void delete(String id) {
- delete(null, id);
- }
+ void delete(StorageAddress address, String id);
/**
- * Deletes data set with given identifier from location defined by StorageOptions
- *
- * @param options storage option to apply during the operation
- * @param id unique identifier of the data set
- */
- void delete(StorageAddress options, String id);
-
- /**
- * Restores previously deleted set with given identifier from the default location in the underlying data store
+ * Restores previously deleted data set with given identifier from location defined by StorageAddress
*
+ * @param address storage address to apply during the operation, can be null to use the default location
* @param id unique identifier of the data set
* @throws ItemNotFoundException in case restore cannot be performed
*/
- default void restore(String id) {
- restore(null, id);
- }
+ void restore(StorageAddress address, String id);
/**
- * Restores previously deleted data set with given identifier from location defined by StorageOptions
+ * Checks if data set with given identifier exists in the location defined by StorageAddress
*
- * @param options storage option to apply during the operation
+ * @param address storage address to apply during the operation, can be null to use the default location
* @param id unique identifier of the data set
- * @throws ItemNotFoundException in case restore cannot be performed
+ * @return true if exists false otherwise
*/
- void restore(StorageAddress options, String id);
+ boolean exists(StorageAddress address, String id);
/**
- * Checks if data set with given identifier exists in the default location in the underlying data store
+ * Stores sub model associated with given data set in location defined by StorageAddress
*
+ * @param address storage address to apply during the operation, can be null to use the default location
* @param id unique identifier of the data set
- * @return true if exists false otherwise
+ * @param kind kind of the sub model e.g. waypoints
+ * @param content content of the sub model to be stored
*/
- default boolean exists(String id) {
- return exists(null, id);
- }
+ void storeSubModel(StorageAddress address, String id, SubModelKind kind, StorageContent content);
/**
- * Checks if data set with given identifier exists in the location defined by StorageOptions
+ * Updates already stored sub model associated with given data set in location defined by
+ * StorageAddress.
+ *
+ * Contrary to {@link #storeSubModel(StorageAddress, String, SubModelKind, StorageContent)}, implementations are
+ * allowed to consider the update as a best effort operation, as updates are usually issued repeatedly while solving.
*
- * @param options storage option to apply during the operation
+ * @param address storage address to apply during the operation, can be null to use the default location
* @param id unique identifier of the data set
- * @param kind type of sub model
- * @return true if exists false otherwise
+ * @param kind kind of the sub model e.g. waypoints
+ * @param content content of the sub model to be stored
*/
- boolean existsSubModel(StorageAddress options, String id, SubModelKind kind);
+ void updateSubModel(StorageAddress address, String id, SubModelKind kind, StorageContent content);
/**
- * Checks if data set with given identifier exists in the default location in the underlying data store
+ * Retrieves content of the sub model associated with data set with given identifier from location defined by
+ * StorageAddress.
+ *
+ * It is the responsibility of the caller to close the returned stream.
*
+ * @param address storage address to apply during the operation, can be null to use the default location
* @param id unique identifier of the data set
- * @param kind type of sub model
- * @return true if exists false otherwise
+ * @param kind kind of the sub model e.g. waypoints
+ * @return stream with the content of the sub model or null if no such sub model is found
*/
- default boolean existsSubModel(String id, SubModelKind kind) {
- return existsSubModel(null, id, kind);
- }
+ InputStream getSubModel(StorageAddress address, String id, SubModelKind kind);
/**
- * Checks if data set with given identifier exists in the location defined by StorageOptions
+ * Checks if sub model of given kind exists for the data set with given identifier in the location defined by
+ * StorageAddress
*
- * @param options storage option to apply during the operation
+ * @param address storage address to apply during the operation, can be null to use the default location
* @param id unique identifier of the data set
+ * @param kind kind of the sub model e.g. waypoints
* @return true if exists false otherwise
*/
- boolean exists(StorageAddress options, String id);
+ boolean existsSubModel(StorageAddress address, String id, SubModelKind kind);
/**
- * Lists data sets (as statuses) stored in the default location of the underlying data store
+ * Tells whether the underlying data store keeps the attributes given to it as part of
+ * {@link StorageContent#attributes()}, so that they can be read back by
+ * {@link #list(StorageAddress, int, int)} without reading the content itself.
*
- * @param pageNumber number of page to return (0-based)
- * @param pageSize number of data sets to return per page
- * @return non null list of found data sets as statuses
+ * @return true if attributes are supported, which is the default
*/
- default List> list(int pageNumber, int pageSize) {
- return list(null, pageNumber, pageSize);
+ default boolean supportsAttributes() {
+ return true;
}
/**
- * Lists data sets (as statuses) stored in the location defined by StorageOptions
+ * Lists data sets stored in the location defined by StorageAddress
*
- * @param options storage option to apply during the operation
+ * @param address storage address to apply during the operation, can be null to use the default location
* @param pageNumber number of page to return (0-based)
* @param pageSize number of data sets to return per page
- * @return non null list of found data sets as statuses
+ * @return non null list of found data sets
*/
- List> list(StorageAddress options, int pageNumber, int pageSize);
+ List list(StorageAddress address, int pageNumber, int pageSize);
/**
- * Retrieves sub resource associated with data set with given identifier in the default location in the underlying data
- * store as a stream of data
+ * Cleans up the storage starting in the location defined by StorageAddress. Removes all resources
+ * regardless of their type e.g. data sets, waypoints etc
*
- * @param id unique identifier of the data set
- * @param subModelKind kind of the sub resource e.g. waypoints
- * @param output output stream where data should be written
+ * @param address storage address to apply during the operation, can be null to use the default location
*/
- default void getSubModelStream(String id, SubModelKind subModelKind, OutputStream output) {
- getSubModelStream(null, id, subModelKind, output);
- }
-
- /**
- * Retrieves sub resource associated with data set with given identifier in the default location in the underlying data
- * store as a stream of data
- *
- * @param options storage option to apply during the operation
- * @param id unique identifier of the data set
- * @param subModelKind kind of the sub resource e.g. waypoints
- * @param output output stream where data should be written
- */
- void getSubModelStream(StorageAddress options, String id, SubModelKind subModelKind, OutputStream output);
-
- /**
- * Retrieves sub resource associated with data set with given identifier in the default location in the underlying data
- * store
- *
- * @param type representing the returned value
- * @param id unique identifier of the data set
- * @param subModelKind kind of the sub resource e.g. waypoints
- * @param clazz class that the sub resource should be unmarshalled to
- * @return loaded sub resource or null if no sub resource found
- */
- default T getSubModel(String id, SubModelKind subModelKind, Class clazz) {
- return getSubModel(null, id, subModelKind, clazz);
- }
-
- /**
- * Retrieves sub resource associated with data set with given identifier in the location defined by
- * StorageOptions
- *
- * @param type representing the returned value
- * @param options storage option to apply during the operation
- * @param id unique identifier of the data set
- * @param subModelKind kind of the sub resource e.g. waypoints
- * @param clazz class that the sub resource should be unmarshalled
- * @return loaded sub resource or null if no sub resource found
- */
- T getSubModel(StorageAddress options, String id, SubModelKind subModelKind, Class clazz);
-
- default T getSubModel(String id, SubModelKind config, TypeReference configurationClass) {
- return getSubModel(null, id, config, configurationClass);
- }
-
- T getSubModel(StorageAddress options, String id, SubModelKind config, TypeReference configurationClass);
-
- /**
- * Stores sub resource associated with given data set given by identifier in the default location of the underlying data
- * store
- *
- * @param id unique identifier of the data set
- * @param subModelKind kind of the sub resource e.g. waypoints
- * @param subModel the sub resource to be stored
- */
- default void storeSubModel(String id, SubModelKind subModelKind, Object subModel) {
- storeSubModel(null, id, subModelKind, subModel);
- }
-
- /**
- * Stores sub resource associated with given data set given by identifier in the location defined by
- * StorageOptions
- *
- * @param options storage option to apply during the operation
- * @param id unique identifier of the data set
- * @param subModelKind kind of the sub resource e.g. waypoints
- * @param subModel the sub resource to be stored
- */
- void storeSubModel(StorageAddress options, String id, SubModelKind subModelKind, Object subModel);
-
- /**
- * Stores a sub resource verbatim from a binary stream in the default location of the underlying data store, without any
- * marshalling. Use this (rather than {@link #storeSubModel}) for binary payloads such as archives, so that the bytes can be
- * retrieved unchanged through {@link #getSubModelStream}.
- *
- * @param id unique identifier of the data set
- * @param subModelKind kind of the sub resource e.g. execution profile artifacts
- * @param input the binary content to be stored; the caller is responsible for closing it
- */
- default void storeSubModelStream(String id, SubModelKind subModelKind, InputStream input) {
- storeSubModelStream(null, id, subModelKind, input);
- }
-
- /**
- * Stores a sub resource verbatim from a binary stream in the location defined by StorageOptions, without any
- * marshalling. Use this (rather than {@link #storeSubModel}) for binary payloads such as archives, so that the bytes can be
- * retrieved unchanged through {@link #getSubModelStream}.
- *
- * @param options storage option to apply during the operation
- * @param id unique identifier of the data set
- * @param subModelKind kind of the sub resource e.g. execution profile artifacts
- * @param input the binary content to be stored; the caller is responsible for closing it
- */
- void storeSubModelStream(StorageAddress options, String id, SubModelKind subModelKind, InputStream input);
-
- /**
- * Stores sub resource associated with given data set given by identifier in the default location of the underlying data
- * store
- *
- * @param id unique identifier of the data set
- * @param subModelKind kind of the sub resource e.g. waypoints
- * @param subModel the sub resource to be stored
- */
- default void updateSubModel(String id, SubModelKind subModelKind, Object subModel) {
- updateSubModel(null, id, subModelKind, subModel);
- }
-
- /**
- * Stores sub resource associated with given data set given by identifier in the location defined by
- * StorageOptions
- *
- * @param options storage option to apply during the operation
- * @param id unique identifier of the data set
- * @param subModelKind kind of the sub resource e.g. waypoints
- * @param subModel the sub resource to be stored
- */
- void updateSubModel(StorageAddress options, String id, SubModelKind subModelKind, Object subModel);
-
- /**
- * Cleans up the storage starting in the default location in the underlying data store. Removes all resources regardless of
- * their type e.g. data sets, waypoints etc
- */
- default void clean() {
- clean(null);
- }
-
- /**
- * Cleans up the storage starting in the location defined by
- * StorageOptions. Removes all resources regardless of
- * their type e.g. data sets, waypoints etc
- *
- * @param options storage option to apply during the operation
- */
- void clean(StorageAddress options);
+ void clean(StorageAddress address);
/**
* Create required data store specific settings to be able to storage data in given location.
@@ -361,6 +184,4 @@ default void clean() {
* @param id named location to be destroyed in the underlying storage e.g. name of the bucket
*/
void destroy(String id);
-
- Class clazz();
}
diff --git a/service/definition/src/main/java/ai/timefold/solver/service/definition/internal/storage/StorageContent.java b/service/definition/src/main/java/ai/timefold/solver/service/definition/internal/storage/StorageContent.java
new file mode 100644
index 00000000000..44eff70e5da
--- /dev/null
+++ b/service/definition/src/main/java/ai/timefold/solver/service/definition/internal/storage/StorageContent.java
@@ -0,0 +1,80 @@
+package ai.timefold.solver.service.definition.internal.storage;
+
+import java.io.ByteArrayInputStream;
+import java.io.InputStream;
+import java.util.Map;
+import java.util.Objects;
+import java.util.concurrent.atomic.AtomicBoolean;
+import java.util.function.Supplier;
+
+/**
+ * Content to be written into a {@link Storage}.
+ *
+ * The content is exposed as a stream so that implementations do not have to keep it in memory. The length is provided
+ * as well, as some data stores require the size of the content up front.
+ *
+ * Content that was handed over as bytes can be read as often as needed, every call to {@link #stream()} opens a new
+ * stream over them. Content that is being forwarded from a stream somebody else opened, such as the body of a request
+ * the access service relays, can only be read once. {@link #repeatable()} tells the two apart, so that a failed
+ * operation is only ever retried when the content can actually be sent again.
+ *
+ * @param source opens a stream with the content, called once per attempt to write it
+ * @param length number of bytes available in the stream, negative when it is not known
+ * @param attributes additional attributes (such as object metadata or tags) to be associated with the content,
+ * never null but possibly empty
+ * @param repeatable whether the content can be read more than once
+ */
+public record StorageContent(Supplier source, long length, Map attributes,
+ boolean repeatable) {
+
+ public StorageContent {
+ Objects.requireNonNull(source, "source cannot be null");
+ if (attributes == null) {
+ attributes = Map.of();
+ }
+ }
+
+ /**
+ * Content of a stream somebody else opened, which can only be read once and therefore cannot be written again
+ * after a failed attempt.
+ *
+ * @param stream stream with the content, it is the responsibility of the storage to consume it
+ * @param length number of bytes available in the stream, negative when it is not known
+ * @param attributes additional attributes to be associated with the content, can be null
+ */
+ public StorageContent(InputStream stream, long length, Map attributes) {
+ this(once(stream), length, attributes, false);
+ }
+
+ /**
+ * Opens the content for reading. A {@link #repeatable()} content returns a new stream on every call, so that an
+ * attempt that failed half way through does not affect the next one.
+ *
+ * @return stream with the content, never null
+ * @throws IllegalStateException when a content that is not repeatable is opened a second time
+ */
+ public InputStream stream() {
+ return source.get();
+ }
+
+ public static StorageContent of(byte[] content) {
+ return of(content, Map.of());
+ }
+
+ public static StorageContent of(byte[] content, Map attributes) {
+ Objects.requireNonNull(content, "content cannot be null");
+ return new StorageContent(() -> new ByteArrayInputStream(content), content.length, attributes, true);
+ }
+
+ private static Supplier once(InputStream stream) {
+ Objects.requireNonNull(stream, "stream cannot be null");
+ var handedOut = new AtomicBoolean();
+ return () -> {
+ if (handedOut.getAndSet(true)) {
+ throw new IllegalStateException("The content is not repeatable and its stream was already handed out;"
+ + " content that has to be read more than once has to be read into memory first.");
+ }
+ return stream;
+ };
+ }
+}
diff --git a/service/definition/src/main/java/ai/timefold/solver/service/definition/internal/storage/StorageItem.java b/service/definition/src/main/java/ai/timefold/solver/service/definition/internal/storage/StorageItem.java
new file mode 100644
index 00000000000..018a7fa1b9a
--- /dev/null
+++ b/service/definition/src/main/java/ai/timefold/solver/service/definition/internal/storage/StorageItem.java
@@ -0,0 +1,25 @@
+package ai.timefold.solver.service.definition.internal.storage;
+
+import java.util.Map;
+import java.util.Objects;
+
+/**
+ * A data set found in a {@link Storage}.
+ *
+ * @param id unique identifier of the data set, without any storage specific prefixes
+ * @param attributes attributes associated with the data set in the underlying data store, never null but possibly
+ * empty when the data store does not support them
+ */
+public record StorageItem(String id, Map attributes) {
+
+ public StorageItem {
+ Objects.requireNonNull(id, "id cannot be null");
+ if (attributes == null) {
+ attributes = Map.of();
+ }
+ }
+
+ public StorageItem(String id) {
+ this(id, Map.of());
+ }
+}
diff --git a/service/definition/src/main/java/ai/timefold/solver/service/definition/internal/storage/SupportedStorages.java b/service/definition/src/main/java/ai/timefold/solver/service/definition/internal/storage/SupportedStorages.java
deleted file mode 100644
index b4dba4bad5e..00000000000
--- a/service/definition/src/main/java/ai/timefold/solver/service/definition/internal/storage/SupportedStorages.java
+++ /dev/null
@@ -1,29 +0,0 @@
-package ai.timefold.solver.service.definition.internal.storage;
-
-public class SupportedStorages {
-
- public static final String STORAGE_TYPE_PROPERTY = "timefold.storage.type";
- public static final String GOOGLE_CLOUD_STORAGE = "googlecloud";
- public static final String AZURE_STORAGE = "azure";
- public static final String INMEMORY_STORAGE = "inmemory";
- public static final String FILESYSTEM_STORAGE = "filesystem";
- public static final String S3_STORAGE = "s3";
-
- public enum Variant {
- S3(S3_STORAGE),
- GoogleCloud(GOOGLE_CLOUD_STORAGE),
- Azure(AZURE_STORAGE),
- InMemory(INMEMORY_STORAGE),
- FileSystem(FILESYSTEM_STORAGE);
-
- private String identifier;
-
- Variant(String identifier) {
- this.identifier = identifier;
- }
-
- public String identifier() {
- return this.identifier;
- }
- }
-}
diff --git a/service/definition/src/test/java/ai/timefold/solver/service/definition/impl/storage/CompressionUtilsTest.java b/service/definition/src/test/java/ai/timefold/solver/service/definition/impl/storage/CompressionUtilsTest.java
new file mode 100644
index 00000000000..8ca25141043
--- /dev/null
+++ b/service/definition/src/test/java/ai/timefold/solver/service/definition/impl/storage/CompressionUtilsTest.java
@@ -0,0 +1,57 @@
+package ai.timefold.solver.service.definition.impl.storage;
+
+import static org.assertj.core.api.Assertions.assertThat;
+
+import java.io.ByteArrayInputStream;
+import java.io.ByteArrayOutputStream;
+import java.io.IOException;
+import java.io.InputStream;
+import java.nio.charset.StandardCharsets;
+import java.util.zip.GZIPInputStream;
+
+import org.junit.jupiter.api.Test;
+
+class CompressionUtilsTest {
+
+ private static final byte[] CONTENT = "{\"id\":\"dataset-1\",\"value\":\"the content of the dataset\"}"
+ .repeat(100)
+ .getBytes(StandardCharsets.UTF_8);
+
+ @Test
+ void compress_producesACompleteGzipStream() throws IOException {
+ byte[] compressed = CompressionUtils.compress(CONTENT);
+
+ assertThat(CompressionUtils.isCompressed(compressed)).isTrue();
+ try (var gzip = new GZIPInputStream(new ByteArrayInputStream(compressed))) {
+ assertThat(gzip.readAllBytes()).isEqualTo(CONTENT);
+ }
+ }
+
+ @Test
+ void compressedContent_uncompressesToTheOriginal() {
+ assertThat(CompressionUtils.uncompress(CompressionUtils.compress(CONTENT))).isEqualTo(CONTENT);
+ }
+
+ @Test
+ void compressedContent_isDecompressedWhenStreamed() throws IOException {
+ try (InputStream content =
+ CompressionUtils.decompressIfNeeded(new ByteArrayInputStream(CompressionUtils.compress(CONTENT)))) {
+ assertThat(content.readAllBytes()).isEqualTo(CONTENT);
+ }
+ }
+
+ @Test
+ void transferredContent_isCompressedOnceAndUncompressesToTheOriginal() throws IOException {
+ var target = new ByteArrayOutputStream();
+ CompressionUtils.transferDataCompressIfNeeded(new ByteArrayInputStream(CONTENT), target);
+
+ var again = new ByteArrayOutputStream();
+ CompressionUtils.transferDataCompressIfNeeded(new ByteArrayInputStream(target.toByteArray()), again);
+
+ assertThat(again.toByteArray()).as("content that is compressed already is passed on as it is")
+ .isEqualTo(target.toByteArray());
+ try (var gzip = new GZIPInputStream(new ByteArrayInputStream(target.toByteArray()))) {
+ assertThat(gzip.readAllBytes()).isEqualTo(CONTENT);
+ }
+ }
+}
diff --git a/service/definition/src/test/java/ai/timefold/solver/service/definition/internal/platform/ServiceAccountTokenTest.java b/service/definition/src/test/java/ai/timefold/solver/service/definition/internal/platform/ServiceAccountTokenTest.java
new file mode 100644
index 00000000000..1cbcde1b0f9
--- /dev/null
+++ b/service/definition/src/test/java/ai/timefold/solver/service/definition/internal/platform/ServiceAccountTokenTest.java
@@ -0,0 +1,102 @@
+package ai.timefold.solver.service.definition.internal.platform;
+
+import static org.assertj.core.api.Assertions.assertThat;
+
+import java.io.IOException;
+import java.nio.file.Files;
+import java.nio.file.Path;
+import java.time.Duration;
+
+import org.junit.jupiter.api.Test;
+import org.junit.jupiter.api.io.TempDir;
+
+class ServiceAccountTokenTest {
+
+ @TempDir
+ Path directory;
+
+ @Test
+ void theTokenIsSentAsABearerCredential() throws IOException {
+ Path file = tokenFile("a.b.c");
+
+ assertThat(new ServiceAccountToken(file.toString()).authorization()).hasValue("Bearer a.b.c");
+ }
+
+ /**
+ * The file ends with a newline in some clusters, which is not part of the token.
+ */
+ @Test
+ void theTokenIsReadWithoutTheWhitespaceAroundIt() throws IOException {
+ Path file = tokenFile("\na.b.c\n");
+
+ assertThat(new ServiceAccountToken(file.toString()).authorization()).hasValue("Bearer a.b.c");
+ }
+
+ /**
+ * Anything that does not run in a cluster has no token mounted, which is not an error: the request is sent
+ * without one and the service on the other end decides what to make of that.
+ */
+ @Test
+ void thereIsNoTokenWhenNoneIsMounted() {
+ ServiceAccountToken token = new ServiceAccountToken(directory.resolve("not-mounted").toString());
+
+ assertThat(token.authorization()).isEmpty();
+ }
+
+ @Test
+ void thereIsNoTokenWhenTheFileIsEmpty() throws IOException {
+ Path file = tokenFile(" \n");
+
+ assertThat(new ServiceAccountToken(file.toString()).authorization()).isEmpty();
+ }
+
+ /**
+ * The kubelet rewrites the file in place when it rotates the token, so holding on to the first one that was read
+ * would mean sending an expired credential for the rest of the life of the pod.
+ */
+ @Test
+ void aRotatedTokenIsPickedUp() throws IOException {
+ Path file = tokenFile("first");
+ ServiceAccountToken token = new ServiceAccountToken(file, Duration.ZERO);
+ assertThat(token.authorization()).hasValue("Bearer first");
+
+ Files.writeString(file, "second");
+
+ assertThat(token.authorization()).hasValue("Bearer second");
+ }
+
+ /**
+ * Reading the file on every request would be a syscall per storage operation, and the token is good for hours.
+ */
+ @Test
+ void theTokenIsNotReadAgainForEveryRequest() throws IOException {
+ Path file = tokenFile("first");
+ ServiceAccountToken token = new ServiceAccountToken(file, Duration.ofHours(1));
+ assertThat(token.authorization()).hasValue("Bearer first");
+
+ Files.writeString(file, "second");
+
+ assertThat(token.authorization()).hasValue("Bearer first");
+ }
+
+ /**
+ * A pod that is not mounted a token may still be given one later on, so the absence of the file is not final
+ * either.
+ */
+ @Test
+ void aTokenThatAppearsLaterIsPickedUp() throws IOException {
+ Path file = directory.resolve("token");
+ ServiceAccountToken token = new ServiceAccountToken(file, Duration.ZERO);
+ assertThat(token.authorization()).isEmpty();
+
+ Files.writeString(file, "mounted");
+
+ assertThat(token.authorization()).hasValue("Bearer mounted");
+ }
+
+ private Path tokenFile(String content) throws IOException {
+ Path file = directory.resolve("token");
+ Files.writeString(file, content);
+ return file;
+ }
+}
diff --git a/service/definition/src/test/java/ai/timefold/solver/service/definition/internal/storage/AbstractStorageServiceRetryTest.java b/service/definition/src/test/java/ai/timefold/solver/service/definition/internal/storage/AbstractStorageServiceRetryTest.java
new file mode 100644
index 00000000000..0ef7eb23aca
--- /dev/null
+++ b/service/definition/src/test/java/ai/timefold/solver/service/definition/internal/storage/AbstractStorageServiceRetryTest.java
@@ -0,0 +1,306 @@
+package ai.timefold.solver.service.definition.internal.storage;
+
+import static org.assertj.core.api.Assertions.assertThat;
+import static org.assertj.core.api.Assertions.assertThatThrownBy;
+
+import java.io.ByteArrayInputStream;
+import java.io.IOException;
+import java.io.InputStream;
+import java.io.UncheckedIOException;
+import java.nio.charset.StandardCharsets;
+import java.util.ArrayList;
+import java.util.List;
+import java.util.Map;
+import java.util.concurrent.atomic.AtomicInteger;
+
+import ai.timefold.solver.service.definition.api.ModelConfigOverrides;
+import ai.timefold.solver.service.definition.api.ModelConstraintJustification;
+import ai.timefold.solver.service.definition.api.ModelInput;
+import ai.timefold.solver.service.definition.api.ModelOutput;
+import ai.timefold.solver.service.definition.api.domain.Configuration;
+import ai.timefold.solver.service.definition.api.metrics.ModelInputMetrics;
+import ai.timefold.solver.service.definition.api.metrics.ModelOutputMetrics;
+import ai.timefold.solver.service.definition.impl.storage.CompressionUtils;
+import ai.timefold.solver.service.definition.impl.storage.StorageObjectMapperWrapper;
+import ai.timefold.solver.service.definition.internal.error.ErrorCodes;
+import ai.timefold.solver.service.definition.internal.error.ItemNotFoundException;
+import ai.timefold.solver.service.definition.internal.error.TimefoldRuntimeException;
+
+import org.junit.jupiter.api.BeforeEach;
+import org.junit.jupiter.api.Test;
+
+import com.fasterxml.jackson.core.type.TypeReference;
+import com.fasterxml.jackson.databind.ObjectMapper;
+
+/**
+ * A data store, or the access service in front of it, can be briefly unavailable. The storage service still holds the
+ * object at that point, so it can serialize it again and send it once more.
+ */
+class AbstractStorageServiceRetryTest {
+
+ private FailingStorage storage;
+
+ private TestStorageService service;
+
+ @BeforeEach
+ void setUp() {
+ storage = new FailingStorage();
+ service = new TestStorageService(storage, new StorageObjectMapperWrapper(new ObjectMapper()));
+ service.setRetry(3, "PT0.01S", "PT0.02S");
+ }
+
+ @Test
+ void aWriteThatFailedRecoverablyIsSentAgainInFull() {
+ storage.failWrites(2, recoverable());
+
+ service.storeModelOutput("run-1", new TestOutput("kept"));
+
+ assertThat(storage.writeAttempts()).isEqualTo(3);
+ assertThat(storage.storedBodies()).as("every attempt sends the whole content")
+ .hasSize(3)
+ .allSatisfy(body -> assertThat(body).isEqualTo(storage.storedBodies().getFirst()))
+ .allSatisfy(body -> assertThat(body).isNotEmpty());
+ }
+
+ @Test
+ void aReadThatFailedRecoverablyIsAttemptedAgain() {
+ storage.content(compressed("{\"value\":\"kept\"}"));
+ storage.failReads(1, recoverable());
+
+ TestOutput output = service.getModelOutput("run-1");
+
+ assertThat(output).isEqualTo(new TestOutput("kept"));
+ assertThat(storage.readAttempts()).isEqualTo(2);
+ }
+
+ @Test
+ void aFailureTheStorageDoesNotCallRecoverableIsNotAttemptedAgain() {
+ storage.failWrites(1, new TimefoldRuntimeException(ErrorCodes.STORAGE_UNABLE_TO_WRITE, "rejected", false));
+
+ assertThatThrownBy(() -> service.storeModelOutput("run-1", new TestOutput("kept")))
+ .isInstanceOf(TimefoldRuntimeException.class)
+ .hasMessage("rejected");
+
+ assertThat(storage.writeAttempts()).isEqualTo(1);
+ }
+
+ @Test
+ void aMissingDataSetIsNotAttemptedAgain() {
+ storage.failReads(1, new ItemNotFoundException(ErrorCodes.STORAGE_NO_JOB_FOUND, "no such dataset"));
+
+ assertThatThrownBy(() -> service.getModelOutput("run-1"))
+ .isInstanceOf(ItemNotFoundException.class);
+
+ assertThat(storage.readAttempts()).isEqualTo(1);
+ }
+
+ @Test
+ void theLastFailureIsReportedOnceTheAttemptsAreUsedUp() {
+ storage.failWrites(Integer.MAX_VALUE, recoverable());
+
+ assertThatThrownBy(() -> service.storeModelOutput("run-1", new TestOutput("kept")))
+ .isInstanceOf(TimefoldRuntimeException.class)
+ .hasMessage("the data store is away");
+
+ assertThat(storage.writeAttempts()).isEqualTo(3);
+ }
+
+ @Test
+ void contentThatCannotBeReadAgainIsNotAttemptedAgain() {
+ StorageContent forwarded = new StorageContent(new ByteArrayInputStream("body".getBytes(StandardCharsets.UTF_8)),
+ 4, Map.of());
+ assertThat(forwarded.repeatable()).isFalse();
+ AtomicInteger attempts = new AtomicInteger();
+
+ assertThatThrownBy(() -> service.writeWithRetry("run-1", "forward the content", forwarded, () -> {
+ attempts.incrementAndGet();
+ throw recoverable();
+ })).isInstanceOf(TimefoldRuntimeException.class);
+
+ assertThat(attempts).hasValue(1);
+ }
+
+ private static TimefoldRuntimeException recoverable() {
+ return new TimefoldRuntimeException(ErrorCodes.STORAGE_UNKNOWN, "the data store is away", true);
+ }
+
+ private static byte[] compressed(String json) {
+ return CompressionUtils.compress(json.getBytes(StandardCharsets.UTF_8));
+ }
+
+ public record TestOutput(String value) implements ModelOutput {
+ }
+
+ private static final class TestStorageService
+ extends
+ AbstractStorageService {
+
+ private TestStorageService(Storage storage, StorageObjectMapperWrapper wrapper) {
+ super(storage, wrapper);
+ }
+
+ @Override
+ protected Class> getModelInputClass() {
+ return ModelInput.class;
+ }
+
+ @Override
+ protected Class> getModelOutputClass() {
+ return TestOutput.class;
+ }
+
+ @Override
+ protected Class> getInputMetricsClass() {
+ return ModelInputMetrics.class;
+ }
+
+ @Override
+ protected Class> getOutputMetricsClass() {
+ return ModelOutputMetrics.class;
+ }
+
+ @Override
+ protected TypeReference> getConfigurationClass() {
+ return new TypeReference<>() {
+ };
+ }
+ }
+
+ /**
+ * Storage that fails a given number of times before it starts working, recording what each attempt handed it.
+ */
+ private static final class FailingStorage implements Storage {
+
+ private final List storedBodies = new ArrayList<>();
+ private final AtomicInteger writeAttempts = new AtomicInteger();
+ private final AtomicInteger readAttempts = new AtomicInteger();
+
+ private int writeFailures;
+ private int readFailures;
+ private RuntimeException writeFailure;
+ private RuntimeException readFailure;
+ private byte[] content;
+
+ void failWrites(int times, RuntimeException failure) {
+ this.writeFailures = times;
+ this.writeFailure = failure;
+ }
+
+ void failReads(int times, RuntimeException failure) {
+ this.readFailures = times;
+ this.readFailure = failure;
+ }
+
+ void content(byte[] content) {
+ this.content = content;
+ }
+
+ int writeAttempts() {
+ return writeAttempts.get();
+ }
+
+ int readAttempts() {
+ return readAttempts.get();
+ }
+
+ List storedBodies() {
+ return storedBodies;
+ }
+
+ @Override
+ public void store(StorageAddress address, String id, StorageContent content) {
+ writeAttempts.incrementAndGet();
+ storedBodies.add(readAll(content));
+ if (writeFailures-- > 0) {
+ throw writeFailure;
+ }
+ }
+
+ @Override
+ public InputStream get(StorageAddress address, String id) {
+ readAttempts.incrementAndGet();
+ if (readFailures-- > 0) {
+ throw readFailure;
+ }
+ return new ByteArrayInputStream(content);
+ }
+
+ private static byte[] readAll(StorageContent content) {
+ try (InputStream stream = content.stream()) {
+ return stream.readAllBytes();
+ } catch (IOException e) {
+ throw new UncheckedIOException(e);
+ }
+ }
+
+ @Override
+ public void update(StorageAddress address, String id, StorageContent content) {
+ store(address, id, content);
+ }
+
+ @Override
+ public void complete(StorageAddress address, String id, StorageContent content) {
+ store(address, id, content);
+ }
+
+ @Override
+ public void delete(StorageAddress address, String id) {
+ // life cycle method not needed for this test
+ }
+
+ @Override
+ public void restore(StorageAddress address, String id) {
+ // life cycle method not needed for this test
+ }
+
+ @Override
+ public boolean exists(StorageAddress address, String id) {
+ return content != null;
+ }
+
+ @Override
+ public void storeSubModel(StorageAddress address, String id, SubModelKind kind, StorageContent content) {
+ store(address, id, content);
+ }
+
+ @Override
+ public void updateSubModel(StorageAddress address, String id, SubModelKind kind, StorageContent content) {
+ store(address, id, content);
+ }
+
+ @Override
+ public InputStream getSubModel(StorageAddress address, String id, SubModelKind kind) {
+ return get(address, id);
+ }
+
+ @Override
+ public boolean existsSubModel(StorageAddress address, String id, SubModelKind kind) {
+ return content != null;
+ }
+
+ @Override
+ public List list(StorageAddress address, int pageNumber, int pageSize) {
+ return List.of();
+ }
+
+ @Override
+ public void clean(StorageAddress address) {
+ // life cycle method not needed for this test
+ }
+
+ @Override
+ public void create(String location, StorageConfiguration configuration) {
+ // life cycle method not needed for this test
+ }
+
+ @Override
+ public void reconfigure(String location, StorageConfiguration configuration) {
+ // life cycle method not needed for this test
+ }
+
+ @Override
+ public void destroy(String id) {
+ // life cycle method not needed for this test
+ }
+ }
+}
diff --git a/service/facade/service-parent/pom.xml b/service/facade/service-parent/pom.xml
index 1769296e1ef..7b0e1b69218 100644
--- a/service/facade/service-parent/pom.xml
+++ b/service/facade/service-parent/pom.xml
@@ -510,22 +510,6 @@
ai.timefold.solver.enterprise
timefold-solver-enterprise-service
-