diff --git a/build/bom/pom.xml b/build/bom/pom.xml index 68452ddd2dd..e6b8da897d3 100644 --- a/build/bom/pom.xml +++ b/build/bom/pom.xml @@ -324,6 +324,17 @@ ${version.ai.timefold.solver} sources + + ai.timefold.solver + timefold-solver-service-storage-inmemory + ${version.ai.timefold.solver} + + + ai.timefold.solver + timefold-solver-service-storage-inmemory + ${version.ai.timefold.solver} + sources + 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 - - ai.timefold.solver.enterprise - timefold-solver-enterprise-service-storage-fs - - - ai.timefold.solver.enterprise - timefold-solver-enterprise-service-storage-azure - - - ai.timefold.solver.enterprise - timefold-solver-enterprise-service-storage-gcs - - - ai.timefold.solver.enterprise - timefold-solver-enterprise-service-storage-s3 - ai.timefold.solver timefold-solver-service-defaults diff --git a/service/maps/service-client/src/main/java/ai/timefold/solver/service/maps/service/client/impl/MapServiceClient.java b/service/maps/service-client/src/main/java/ai/timefold/solver/service/maps/service/client/impl/MapServiceClient.java index bce5c2dc192..a98d2db2bd6 100644 --- a/service/maps/service-client/src/main/java/ai/timefold/solver/service/maps/service/client/impl/MapServiceClient.java +++ b/service/maps/service-client/src/main/java/ai/timefold/solver/service/maps/service/client/impl/MapServiceClient.java @@ -5,11 +5,13 @@ import ai.timefold.solver.service.maps.service.integration.internal.MapServiceApi; import ai.timefold.solver.service.maps.service.integration.internal.MapServiceHealthCheckApi; +import org.eclipse.microprofile.rest.client.annotation.RegisterClientHeaders; import org.eclipse.microprofile.rest.client.annotation.RegisterProvider; import org.eclipse.microprofile.rest.client.inject.RegisterRestClient; @RegisterRestClient(configKey = "map-service") @RegisterProvider(MapServiceExceptionMapper.class) +@RegisterClientHeaders(ServiceAccountTokenHeadersFactory.class) public interface MapServiceClient extends MapServiceApi, MapServiceHealthCheckApi, MapManagementApi { } diff --git a/service/maps/service-client/src/main/java/ai/timefold/solver/service/maps/service/client/impl/ServiceAccountTokenHeadersFactory.java b/service/maps/service-client/src/main/java/ai/timefold/solver/service/maps/service/client/impl/ServiceAccountTokenHeadersFactory.java new file mode 100644 index 00000000000..d6a758ab3b5 --- /dev/null +++ b/service/maps/service-client/src/main/java/ai/timefold/solver/service/maps/service/client/impl/ServiceAccountTokenHeadersFactory.java @@ -0,0 +1,39 @@ +package ai.timefold.solver.service.maps.service.client.impl; + +import jakarta.enterprise.context.ApplicationScoped; +import jakarta.inject.Inject; +import jakarta.ws.rs.core.MultivaluedMap; + +import ai.timefold.solver.service.definition.internal.platform.ServiceAccountToken; + +import org.eclipse.microprofile.config.inject.ConfigProperty; +import org.eclipse.microprofile.rest.client.ext.ClientHeadersFactory; + +/** + * Sends the service account token of the pod along with every call to the map service, so that a deployment which + * puts the access service in front of the map service can tell which solver worker is calling. + *

+ * A model that does not run in a cluster has no token mounted and calls the map service without one, the same way it + * did before. + */ +@ApplicationScoped +public class ServiceAccountTokenHeadersFactory implements ClientHeadersFactory { + + private final ServiceAccountToken token; + + @Inject + public ServiceAccountTokenHeadersFactory( + @ConfigProperty(name = ServiceAccountToken.TOKEN_FILE_PROPERTY, + defaultValue = ServiceAccountToken.DEFAULT_TOKEN_FILE) String tokenFile) { + this.token = new ServiceAccountToken(tokenFile); + } + + @Override + public MultivaluedMap update(MultivaluedMap incomingHeaders, + MultivaluedMap outgoingHeaders) { + token.authorization() + .ifPresent(authorization -> outgoingHeaders.putSingle(ServiceAccountToken.AUTHORIZATION_HEADER, + authorization)); + return outgoingHeaders; + } +} diff --git a/service/maps/service-client/src/test/java/ai/timefold/solver/service/maps/service/client/impl/ServiceAccountTokenHeadersFactoryTest.java b/service/maps/service-client/src/test/java/ai/timefold/solver/service/maps/service/client/impl/ServiceAccountTokenHeadersFactoryTest.java new file mode 100644 index 00000000000..db8312ef3a2 --- /dev/null +++ b/service/maps/service-client/src/test/java/ai/timefold/solver/service/maps/service/client/impl/ServiceAccountTokenHeadersFactoryTest.java @@ -0,0 +1,42 @@ +package ai.timefold.solver.service.maps.service.client.impl; + +import static org.assertj.core.api.Assertions.assertThat; + +import java.io.IOException; +import java.nio.file.Files; +import java.nio.file.Path; + +import jakarta.ws.rs.core.MultivaluedHashMap; +import jakarta.ws.rs.core.MultivaluedMap; + +import org.junit.jupiter.api.Test; +import org.junit.jupiter.api.io.TempDir; + +class ServiceAccountTokenHeadersFactoryTest { + + @TempDir + Path directory; + + @Test + void theServiceAccountTokenIsSentWithTheCall() throws IOException { + Path file = directory.resolve("token"); + Files.writeString(file, "a.b.c"); + + MultivaluedMap outgoing = update(file.toString()); + + assertThat(outgoing.getFirst("Authorization")).isEqualTo("Bearer a.b.c"); + } + + @Test + void aModelWithoutAMountedTokenCallsTheMapServiceWithoutOne() { + MultivaluedMap outgoing = update(directory.resolve("not-mounted").toString()); + + assertThat(outgoing).isEmpty(); + } + + private static MultivaluedMap update(String tokenFile) { + MultivaluedMap outgoing = new MultivaluedHashMap<>(); + new ServiceAccountTokenHeadersFactory(tokenFile).update(new MultivaluedHashMap<>(), outgoing); + return outgoing; + } +} diff --git a/service/maps/service-client/src/test/java/ai/timefold/solver/service/maps/service/client/util/DummyStorageService.java b/service/maps/service-client/src/test/java/ai/timefold/solver/service/maps/service/client/util/DummyStorageService.java index 17ab787a3e5..38503fe844c 100644 --- a/service/maps/service-client/src/test/java/ai/timefold/solver/service/maps/service/client/util/DummyStorageService.java +++ b/service/maps/service-client/src/test/java/ai/timefold/solver/service/maps/service/client/util/DummyStorageService.java @@ -21,6 +21,11 @@ protected Class getModelInputClass() { return null; } + @Override + protected Class getModelOutputClass() { + return null; + } + @Override protected Class getInputMetricsClass() { return null; diff --git a/service/pom.xml b/service/pom.xml index 85774b89100..470b3ced379 100644 --- a/service/pom.xml +++ b/service/pom.xml @@ -28,6 +28,7 @@ build/build-support build/config definition + storage-inmemory json jackson rest diff --git a/service/quarkus/deployment/pom.xml b/service/quarkus/deployment/pom.xml index 5f4449120c9..2d4a583dd16 100644 --- a/service/quarkus/deployment/pom.xml +++ b/service/quarkus/deployment/pom.xml @@ -24,6 +24,11 @@ ai.timefold.solver timefold-solver-service-definition + + ai.timefold.solver + timefold-solver-service-storage-inmemory + test + ai.timefold.solver timefold-solver-service-rest diff --git a/service/quarkus/deployment/src/main/java/ai/timefold/solver/service/quarkus/deployment/TimefoldStorageProcessor.java b/service/quarkus/deployment/src/main/java/ai/timefold/solver/service/quarkus/deployment/TimefoldStorageProcessor.java index 9427aaed5d6..aff405d3bb8 100644 --- a/service/quarkus/deployment/src/main/java/ai/timefold/solver/service/quarkus/deployment/TimefoldStorageProcessor.java +++ b/service/quarkus/deployment/src/main/java/ai/timefold/solver/service/quarkus/deployment/TimefoldStorageProcessor.java @@ -7,21 +7,14 @@ import java.util.Collection; import java.util.List; -import jakarta.annotation.PostConstruct; -import jakarta.annotation.PreDestroy; import jakarta.enterprise.context.ApplicationScoped; -import jakarta.inject.Inject; import ai.timefold.solver.service.definition.api.domain.Configuration; -import ai.timefold.solver.service.definition.impl.storage.inmemory.InMemoryStorage; import ai.timefold.solver.service.definition.internal.storage.AbstractStorageService; -import ai.timefold.solver.service.definition.internal.storage.Storage; -import ai.timefold.solver.service.definition.internal.storage.SupportedStorages; import ai.timefold.solver.service.quarkus.deployment.builditem.ModelComponentsBuildItem; import org.eclipse.microprofile.config.inject.ConfigProperty; import org.jboss.jandex.AnnotationInstance; -import org.jboss.jandex.AnnotationInstanceBuilder; import org.jboss.jandex.AnnotationValue; import org.jboss.jandex.ClassInfo; import org.jboss.jandex.DotName; @@ -32,7 +25,6 @@ import io.quarkus.arc.deployment.GeneratedBeanBuildItem; import io.quarkus.arc.deployment.GeneratedBeanGizmoAdaptor; -import io.quarkus.arc.lookup.LookupIfProperty; import io.quarkus.deployment.annotations.BuildProducer; import io.quarkus.deployment.annotations.BuildStep; import io.quarkus.deployment.builditem.CombinedIndexBuildItem; @@ -47,8 +39,6 @@ class TimefoldStorageProcessor { - public static final DotName STORAGE = DotName.createSimple(Storage.class.getName()); - public static final DotName STORAGE_SERVICE = DotName.createSimple(AbstractStorageService.class.getName()); public static final DotName ENTERPRISE_STORAGE_SERVICE = @@ -57,154 +47,6 @@ class TimefoldStorageProcessor { private static final String GENERATED_PACKAGE = "ai.timefold.platform.generated.storage."; - /** - * Generating concrete implementation of ai.timefold.solver.service.api.storage.Storage based on - * project dependency with underlying object store such as S3, Google Cloud Storage or Azure BlobStore. - *

- * Looks up what implements the ai.timefold.solver.service.api.ModelOutput and uses it as the actual type of - * data - * to be stored. - * - * @param combinedIndex - index that is used to find types implementing interfaces - * @param generatedClasses - producer to push generated classes - */ - @BuildStep - void generateStorageImpl(CombinedIndexBuildItem combinedIndex, - ModelComponentsBuildItem modelComponentsBuildItem, - BuildProducer generatedClasses) { - - Collection storageImpl = combinedIndex.getIndex().getAllKnownImplementations(STORAGE); - - for (ClassInfo storageClass : storageImpl) { - - ClassInfo modelOutput = modelComponentsBuildItem.getModelOutput(); - GeneratedBeanGizmoAdaptor classOutput = new GeneratedBeanGizmoAdaptor(generatedClasses); - - String generatedName = - GENERATED_PACKAGE + modelOutput.simpleName() + "_" + storageClass.simpleName(); - // create class definition that extends the concrete implementation and implements storage with user model as param - ClassCreator beanCreator = ClassCreator.builder().classOutput(classOutput).className(generatedName) - .signature(SignatureBuilder.forClass().setSuperClass(Type.classType(storageClass.name())) - .addInterface(Type.parameterizedType(Type.classType(Storage.class.getCanonicalName()), - Type.classType(modelOutput.name().toString())))) - .build(); - beanCreator.addAnnotation(ApplicationScoped.class); - - String storageName = storageClass.simpleName() - .substring(0, storageClass.simpleName().indexOf(Storage.class.getSimpleName())).toLowerCase(); - - AnnotationInstanceBuilder builder = AnnotationInstance.builder(LookupIfProperty.class) - .add("name", SupportedStorages.STORAGE_TYPE_PROPERTY).add("stringValue", storageName); - - if (storageClass.name().toString().equals(InMemoryStorage.class.getCanonicalName())) { - builder.add("lookupIfMissing", true); - } - - beanCreator.addAnnotation(builder.build()); - List constructors = storageClass.constructors(); - List parameterTypes = new ArrayList<>(); - if (constructors.size() > 2) { - throw new IllegalStateException("The Storage class must have at most two constructors: " + - "a mandatory one without any parameters for recording, " + - "an optional one with parameters for injection."); - } - if (constructors.getFirst().parametersCount() != 0 && - constructors.getLast().parametersCount() != 0) { - throw new IllegalStateException("The Storage class is missing the mandatory constructor with no parameters."); - } - - for (MethodInfo constructorInfo : constructors) { - parameterTypes.clear(); - - for (MethodParameterInfo param : constructorInfo.parameters()) { - parameterTypes.add(DescriptorUtils.typeToString(param.type())); - } - - // create constructor with matching parameters of the super class - MethodCreator constructor = - beanCreator.getMethodCreator( - MethodDescriptor.ofConstructor(beanCreator.getSuperClass(), - parameterTypes.toArray(String[]::new))); - - ResultHandle thisObj = constructor.getThis(); - - ResultHandle[] params = new ResultHandle[parameterTypes.size()]; - - for (int i = 0; i < params.length; i++) { - params[i] = constructor.getMethodParam(i); - } - - // Invoke Object's constructor - constructor.invokeSpecialMethod( - MethodDescriptor.ofConstructor(beanCreator.getSuperClass(), parameterTypes.toArray(String[]::new)), - thisObj, - params); - - if (constructors.size() == 1 || !parameterTypes.isEmpty()) { - constructor.addAnnotation(Inject.class); - } - - if (Modifier.isPrivate(constructorInfo.flags()) || - !(Modifier.isProtected(constructorInfo.flags()) || - Modifier.isPublic(constructorInfo.flags()))) { - throw new IllegalStateException( - "Storage Constructor cannot be private or package-private; use protected or public."); - } - - // annotate any of the parameters with config property - int index = 0; - for (MethodParameterInfo param : constructorInfo.parameters()) { - - Collection configAnnotations = param.annotations(ConfigProperty.class); - - if (!configAnnotations.isEmpty()) { - - for (AnnotationInstance annotation : configAnnotations) { - AnnotationCreator an = - constructor.getParameterAnnotations(index).addAnnotation(annotation.name().toString()); - - for (AnnotationValue value : annotation.values()) { - an.add(value.name(), value.value()); - } - } - } - index++; - } - constructor.returnValue(thisObj); - } - - // add implementation of clazz method to return the class of the user model - MethodCreator clazzMethod = beanCreator - .getMethodCreator("clazz", Class.class) - .setModifiers(ACC_PROTECTED); - - clazzMethod.returnValue(clazzMethod.loadClassFromTCCL(modelOutput)); - - // override postconstruct and predestroy methods if exists so they are invoked by CDI - List methods = storageClass.methods(); - - for (MethodInfo method : methods) { - - if (method.hasAnnotation(PostConstruct.class)) { - MethodCreator postConstructMethod = beanCreator.getMethodCreator(method.name(), "V"); - postConstructMethod.addAnnotation(PostConstruct.class); - postConstructMethod.invokeSpecialMethod(MethodDescriptor.of(method), postConstructMethod.getThis()); - postConstructMethod.returnVoid(); - } - - if (method.hasAnnotation(PreDestroy.class)) { - MethodCreator preDestroyMethod = beanCreator.getMethodCreator(method.name(), "V"); - preDestroyMethod.addAnnotation(PreDestroy.class); - preDestroyMethod.invokeSpecialMethod(MethodDescriptor.of(method), preDestroyMethod.getThis()); - preDestroyMethod.returnVoid(); - } - } - - beanCreator.close(); - } - - } - /** * Generate concrete implementation of * ai.timefold.solver.service.api.storage.AbstractStorageService @@ -321,6 +163,13 @@ void generateStorageServiceImpl(CombinedIndexBuildItem combinedIndex, getModelInputClassMethod.returnValue(getModelInputClassMethod.loadClassFromTCCL(modelInput)); + // add implementation of getModelOutputClass method to return the class of the user model + MethodCreator getModelOutputClassMethod = beanCreator + .getMethodCreator("getModelOutputClass", Class.class) + .setModifiers(ACC_PROTECTED); + + getModelOutputClassMethod.returnValue(getModelOutputClassMethod.loadClassFromTCCL(modelOutput)); + // add implementation of getInputMetricsClass method to return the class of the user model MethodCreator getInputMetricsClassMethod = beanCreator .getMethodCreator("getInputMetricsClass", Class.class) diff --git a/service/quarkus/deployment/src/test/java/ai/timefold/solver/service/quarkus/deployment/ModelsExtensionStorageGenerationTest.java b/service/quarkus/deployment/src/test/java/ai/timefold/solver/service/quarkus/deployment/ModelsExtensionStorageGenerationTest.java index beee4ae44dc..82e5dc83162 100644 --- a/service/quarkus/deployment/src/test/java/ai/timefold/solver/service/quarkus/deployment/ModelsExtensionStorageGenerationTest.java +++ b/service/quarkus/deployment/src/test/java/ai/timefold/solver/service/quarkus/deployment/ModelsExtensionStorageGenerationTest.java @@ -33,7 +33,7 @@ public class ModelsExtensionStorageGenerationTest { TestdataSolution.class, TestdataConstraintProvider.class, TestdataRest.class, TestdataModelConvertor.class); @Inject - Storage storage; + Storage storage; @Inject AbstractStorageService storageService; diff --git a/service/quarkus/integration-tests/src/test/java/ai/timefold/solver/service/quarkus/deployment/it/ModelsExtensionTest.java b/service/quarkus/integration-tests/src/test/java/ai/timefold/solver/service/quarkus/deployment/it/ModelsExtensionTest.java index 299979e5d67..5a72c75d403 100644 --- a/service/quarkus/integration-tests/src/test/java/ai/timefold/solver/service/quarkus/deployment/it/ModelsExtensionTest.java +++ b/service/quarkus/integration-tests/src/test/java/ai/timefold/solver/service/quarkus/deployment/it/ModelsExtensionTest.java @@ -6,9 +6,9 @@ import jakarta.inject.Inject; import ai.timefold.solver.core.api.score.HardMediumSoftScore; -import ai.timefold.solver.service.definition.impl.storage.inmemory.InMemoryStorage; import ai.timefold.solver.service.definition.internal.storage.AbstractStorageService; import ai.timefold.solver.service.definition.internal.storage.Storage; +import ai.timefold.solver.service.storage.inmemory.InMemoryStorage; import org.junit.jupiter.api.Test; @@ -18,7 +18,7 @@ public class ModelsExtensionTest { @Inject - Storage storage; + Storage storage; @Inject AbstractStorageService storageService; diff --git a/service/quarkus/runtime/pom.xml b/service/quarkus/runtime/pom.xml index 30f798eaf75..312dca96a4e 100644 --- a/service/quarkus/runtime/pom.xml +++ b/service/quarkus/runtime/pom.xml @@ -40,6 +40,10 @@ ai.timefold.solver timefold-solver-service-definition + + ai.timefold.solver + timefold-solver-service-storage-inmemory + ai.timefold.solver timefold-solver-service-rest diff --git a/service/storage-inmemory/pom.xml b/service/storage-inmemory/pom.xml new file mode 100644 index 00000000000..a295bd68d7d --- /dev/null +++ b/service/storage-inmemory/pom.xml @@ -0,0 +1,39 @@ + + 4.0.0 + + ai.timefold.solver + timefold-solver-service-internal-parent + ${revision} + ../pom.xml + + + timefold-solver-service-storage-inmemory + (Preview) Timefold Solver Service Storage: In Memory + + The storage that keeps the data sets in memory, used when no other storage is configured. + It is a module of its own so that the type definitions do not have to depend on ArC, which the storage needs to + declare itself a bean that is only looked up for its own storage type. + This module is in a preview state and thus is a subject to changes. + + + + + ai.timefold.solver + timefold-solver-service-definition + + + io.quarkus + quarkus-arc + + + + + + + io.smallrye + jandex-maven-plugin + + + + diff --git a/service/storage-inmemory/src/main/java/ai/timefold/solver/service/storage/inmemory/InMemoryStorage.java b/service/storage-inmemory/src/main/java/ai/timefold/solver/service/storage/inmemory/InMemoryStorage.java new file mode 100644 index 00000000000..900281b4cdd --- /dev/null +++ b/service/storage-inmemory/src/main/java/ai/timefold/solver/service/storage/inmemory/InMemoryStorage.java @@ -0,0 +1,181 @@ +package ai.timefold.solver.service.storage.inmemory; + +import java.io.ByteArrayInputStream; +import java.io.IOException; +import java.io.InputStream; +import java.io.UncheckedIOException; +import java.util.ArrayList; +import java.util.List; +import java.util.Map; +import java.util.concurrent.ConcurrentHashMap; + +import jakarta.enterprise.context.ApplicationScoped; + +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.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.StorageContent; +import ai.timefold.solver.service.definition.internal.storage.StorageItem; +import ai.timefold.solver.service.definition.internal.storage.SubModelKind; + +import io.quarkus.arc.lookup.LookupIfProperty; + +/** + * Storage that keeps the content in memory, intended for development and testing purposes. + *

+ * The content is kept as it was handed over by the storage service, meaning compressed, so that this storage behaves + * the same way as the ones backed by an actual data store. + */ +@ApplicationScoped +@LookupIfProperty(name = "timefold.storage.type", stringValue = "inmemory", + lookupIfMissing = true) +public class InMemoryStorage implements Storage { + + private static final String DELETED_SUFFIX = ".deleted"; + + private final Map datasets = new ConcurrentHashMap<>(); + + private final Map subModels = new ConcurrentHashMap<>(); + + private final Map> subModelAttributes = new ConcurrentHashMap<>(); + + @Override + public void store(StorageAddress address, String id, StorageContent content) { + datasets.put(id, readAllBytes(content)); + } + + @Override + public void update(StorageAddress address, String id, StorageContent content) { + datasets.put(id, readAllBytes(content)); + } + + @Override + public void complete(StorageAddress address, String id, StorageContent content) { + datasets.put(id, readAllBytes(content)); + } + + @Override + public InputStream get(StorageAddress address, String id) { + byte[] content = datasets.get(id); + if (content == null) { + throw new ItemNotFoundException(ErrorCodes.STORAGE_NO_JOB_FOUND, "Unable to find dataset for id " + id); + } + return new ByteArrayInputStream(content); + } + + @Override + public void delete(StorageAddress address, String id) { + for (SubModelKind kind : SubModelKind.values()) { + byte[] removed = subModels.remove(key(id, kind)); + Map removedAttributes = subModelAttributes.remove(key(id, kind)); + + if (removed != null) { + subModels.put(key(id, kind) + DELETED_SUFFIX, removed); + if (removedAttributes != null) { + subModelAttributes.put(key(id, kind) + DELETED_SUFFIX, removedAttributes); + } + } + } + byte[] dataset = datasets.remove(id); + if (dataset != null) { + datasets.put(id + DELETED_SUFFIX, dataset); + } + } + + @Override + public void restore(StorageAddress address, String id) { + if (!datasets.containsKey(id + DELETED_SUFFIX)) { + throw new ItemNotFoundException(ErrorCodes.STORAGE_NO_JOB_FOUND, + "Run with id " + id + " cannot be restored as it does not exist"); + } + datasets.put(id, datasets.remove(id + DELETED_SUFFIX)); + for (SubModelKind kind : SubModelKind.values()) { + byte[] restored = subModels.remove(key(id, kind) + DELETED_SUFFIX); + Map restoredAttributes = subModelAttributes.remove(key(id, kind) + DELETED_SUFFIX); + + if (restored != null) { + subModels.put(key(id, kind), restored); + if (restoredAttributes != null) { + subModelAttributes.put(key(id, kind), restoredAttributes); + } + } + } + } + + @Override + public boolean exists(StorageAddress address, String id) { + return datasets.containsKey(id); + } + + @Override + public void storeSubModel(StorageAddress address, String id, SubModelKind kind, StorageContent content) { + subModels.put(key(id, kind), readAllBytes(content)); + subModelAttributes.put(key(id, kind), content.attributes()); + } + + @Override + public void updateSubModel(StorageAddress address, String id, SubModelKind kind, StorageContent content) { + storeSubModel(address, id, kind, content); + } + + @Override + public InputStream getSubModel(StorageAddress address, String id, SubModelKind kind) { + byte[] content = subModels.get(key(id, kind)); + + return content == null ? null : new ByteArrayInputStream(content); + } + + @Override + public boolean existsSubModel(StorageAddress address, String id, SubModelKind kind) { + return subModels.containsKey(key(id, kind)); + } + + @Override + public List list(StorageAddress address, int pageNumber, int pageSize) { + String suffix = "_" + SubModelKind.METADATA.id(); + List items = new ArrayList<>(); + for (String key : subModels.keySet()) { + if (key.endsWith(suffix)) { + String id = key.substring(0, key.length() - suffix.length()); + items.add(new StorageItem(id, subModelAttributes.getOrDefault(key, Map.of()))); + } + } + return items.stream().skip((long) pageNumber * pageSize).limit(pageSize).toList(); + } + + @Override + public void clean(StorageAddress address) { + datasets.clear(); + subModels.clear(); + subModelAttributes.clear(); + } + + @Override + public void create(String location, StorageConfiguration configuration) { + // in memory storage does not need to create anything, it is always ready to use + } + + @Override + public void reconfigure(String location, StorageConfiguration configuration) { + // in memory storage does not need to reconfigure anything, it is always ready to use + } + + @Override + public void destroy(String id) { + // in memory storage does not need to destroy anything, it is always ready to use + } + + private static String key(String id, SubModelKind kind) { + return id + "_" + kind.id(); + } + + private static byte[] readAllBytes(StorageContent content) { + try (InputStream stream = content.stream()) { + return stream.readAllBytes(); + } catch (IOException e) { + throw new UncheckedIOException("Unable to read the content to be stored", e); + } + } +} diff --git a/service/worker/pom.xml b/service/worker/pom.xml index 7df1b1f1c7a..f0ae2ae2363 100644 --- a/service/worker/pom.xml +++ b/service/worker/pom.xml @@ -19,6 +19,11 @@ ai.timefold.solver timefold-solver-service-definition + + ai.timefold.solver + timefold-solver-service-storage-inmemory + test + ai.timefold.solver timefold-solver-quarkus-jackson diff --git a/service/worker/src/test/java/ai/timefold/solver/service/worker/impl/DefaultSolverWorkerFacadeTest.java b/service/worker/src/test/java/ai/timefold/solver/service/worker/impl/DefaultSolverWorkerFacadeTest.java index d36835f3baa..958067fd8ad 100644 --- a/service/worker/src/test/java/ai/timefold/solver/service/worker/impl/DefaultSolverWorkerFacadeTest.java +++ b/service/worker/src/test/java/ai/timefold/solver/service/worker/impl/DefaultSolverWorkerFacadeTest.java @@ -12,6 +12,7 @@ import ai.timefold.solver.service.definition.api.domain.Metadata; import ai.timefold.solver.service.definition.api.domain.ModelInputPatchRequest; import ai.timefold.solver.service.definition.api.rest.DatasetSelector; +import ai.timefold.solver.service.definition.impl.storage.StorageObjectMapperWrapper; import ai.timefold.solver.service.definition.internal.error.ItemNotFoundException; import ai.timefold.solver.service.definition.internal.events.DatasetCreatedEvent; import ai.timefold.solver.service.definition.internal.events.DatasetValidateComputeCommand; @@ -40,8 +41,9 @@ class DefaultSolverWorkerFacadeTest { @BeforeEach void setUp() { - var mapper = new ObjectMapper(); - storageService = new TestdataStorageService(new TestdataStorage(mapper)); + // register the Java time module, so that the metadata timestamps can be (de)serialized by the storage + var mapper = new ObjectMapper().findAndRegisterModules(); + storageService = new TestdataStorageService(new TestdataStorage(), new StorageObjectMapperWrapper(mapper)); datasetCreatedEmitter = new RecordingEmitter<>(); validateComputeEmitter = new RecordingEmitter<>(); diff --git a/service/worker/src/test/java/ai/timefold/solver/service/worker/impl/testdata/TestdataStorage.java b/service/worker/src/test/java/ai/timefold/solver/service/worker/impl/testdata/TestdataStorage.java index 4639fbbd8ab..a2801cf5cc8 100644 --- a/service/worker/src/test/java/ai/timefold/solver/service/worker/impl/testdata/TestdataStorage.java +++ b/service/worker/src/test/java/ai/timefold/solver/service/worker/impl/testdata/TestdataStorage.java @@ -1,21 +1,10 @@ package ai.timefold.solver.service.worker.impl.testdata; -import ai.timefold.solver.service.definition.impl.storage.inmemory.InMemoryStorage; - -import com.fasterxml.jackson.databind.ObjectMapper; +import ai.timefold.solver.service.storage.inmemory.InMemoryStorage; /** * Real, in-memory-backed {@code Storage} used by tests instead of a mock, so that reads reflect what was * actually written rather than a stubbed expectation. */ -public class TestdataStorage extends InMemoryStorage { - - public TestdataStorage(ObjectMapper mapper) { - super(mapper); - } - - @Override - public Class clazz() { - return TestdataModelOutput.class; - } -} \ No newline at end of file +public class TestdataStorage extends InMemoryStorage { +} diff --git a/service/worker/src/test/java/ai/timefold/solver/service/worker/impl/testdata/TestdataStorageService.java b/service/worker/src/test/java/ai/timefold/solver/service/worker/impl/testdata/TestdataStorageService.java index 45902595d87..9cd74c58f0a 100644 --- a/service/worker/src/test/java/ai/timefold/solver/service/worker/impl/testdata/TestdataStorageService.java +++ b/service/worker/src/test/java/ai/timefold/solver/service/worker/impl/testdata/TestdataStorageService.java @@ -2,6 +2,7 @@ import ai.timefold.solver.core.api.score.HardSoftScore; import ai.timefold.solver.service.definition.api.domain.Configuration; +import ai.timefold.solver.service.definition.impl.storage.StorageObjectMapperWrapper; import ai.timefold.solver.service.definition.internal.storage.AbstractStorageService; import ai.timefold.solver.service.definition.internal.storage.Storage; @@ -13,8 +14,8 @@ public class TestdataStorageService extends AbstractStorageService { - public TestdataStorageService(Storage storage) { - super(storage); + public TestdataStorageService(Storage storage, StorageObjectMapperWrapper storageObjectMapperWrapper) { + super(storage, storageObjectMapperWrapper); } @Override @@ -22,6 +23,11 @@ protected Class getModelInputClass() { return TestdataModelInput.class; } + @Override + protected Class getModelOutputClass() { + return TestdataModelOutput.class; + } + @Override protected Class getInputMetricsClass() { return TestdataModelInputMetrics.class; @@ -37,4 +43,4 @@ protected TypeReference> getConfigur return new TypeReference<>() { }; } -} \ No newline at end of file +}