Logs
LogStorage and Log are the Lakestream APIs for appending to and reading a log, covering entries, offsets, fencing and buffer ownership.
LogStorage (io.lakestream.api, extends Closeable) is "Level 1": addressable operations on individual logs identified by LogId. Log (extends AutoCloseable) is the per-log managed handle obtained through StreamCatalog.openLog — it wraps a LogStorage, adds a shared entry-index cache, cursor management, retention and search helpers, and fencing. Both interfaces declare that implementations must be safe for concurrent use.
Types documented on this page (13): LogStorage, Log, LogId, LogEntry, LogEntryHeader, EntryHeader, LogEntryIndex, EntryIndex (with nested HeaderWithIndex and IndexType), LogOffset, Position (with nested FileType), FileInfo, LogState, LogStateManager. Log's cursor-management methods (openCursor, openEphemeralCursor, loadCursor, loadAllCursors, deleteCursor) are documented on Cursors instead, alongside LogCursor.
LogStorage
CompletableFuture<LogEntryHeader> append(LogId logId, int numberOfRecords, ByteBuf data);
CompletableFuture<List<LogEntry>> readEntries(
LogId logId, long startOffset, int maxMessageCount, long maxSizeBytes);
CompletableFuture<LogOffset> getFirstOffset(LogId logId);
CompletableFuture<LogOffset> getLastOffset(LogId logId);
CompletableFuture<Long> softTrim(LogId logId, long offsetIncluded);
CompletableFuture<List<EntryIndex>> readIndexRange(LogId logId, long startOffset, long endOffset);
CompletableFuture<Void> hardTrim(LogId logId, long offsetExcluded);
CompletableFuture<Void> deleteLog(LogId logId);
CompletableFuture<List<LogEntry>> readEntriesByIndex(LogId logId, List<EntryIndex> indices,
long startOffset, long maxOffset, int maxMessageCount, long maxSizeBytes,
Predicate<Long> offsetDeleted, Predicate<Long> skipCondition); // default
CompletableFuture<LogOffset> getFirstOffset(LogId logId, boolean includeTrimmed); // default
void preFetchEntries(LogId logId, List<Position> positions); // defaultLogStorage has no create operation: a StreamCatalog implementation allocates log IDs, and callers get them from a stream's committed StreamLayout. deleteLog must be idempotent — deleting a log that is already absent must complete successfully, because lifecycle cleanup can crash after the physical deletion but before recording that it happened, and recovery then replays the same deletion. Deleting a log through a catalog-opened Log handle instead (log.logStorage().deleteLog(...)) is a distinct call path from StreamCatalog.dropStream; Ursa 1.0 fails it — see Implementation status.
LogStorage.getFirstOffset and LogStorage.getLastOffset complete their returned future exceptionally when the log is empty, rather than completing normally with a sentinel value — the Javadoc for both says "or an exceptional future if the log is empty." Log's same-named methods carry no such clause. LogOffset.NOT_FOUND is a sentinel for other uses; it is not what these two LogStorage methods return for an empty log. See Implementation status for a documented deviation.
Log
LogId id();
CompletableFuture<LogEntryHeader> append(int numberOfRecords, ByteBuf data);
CompletableFuture<List<LogEntry>> readEntries(long startOffset, int maxMessageCount, long maxSizeBytes);
CompletableFuture<List<LogEntry>> readEntries(
long startOffset, int maxMessageCount, long maxSizeBytes, boolean includeTrimmed);
CompletableFuture<LogEntry> readEntry(long offset);
CompletableFuture<LogEntryHeader> getEntryMetadata(long offset);
CompletableFuture<EntryIndex> getEntryIndex(long offset);
CompletableFuture<List<EntryIndex>> readIndexRange(long startOffset, long endOffset);
CompletableFuture<List<LogEntryHeader>> getEntryMetadataRange(long startOffset, long endOffset);
CompletableFuture<LogOffset> getFirstOffset();
CompletableFuture<LogOffset> getFirstOffset(boolean includeTrimmed);
CompletableFuture<LogOffset> getLastOffset();
CompletableFuture<Long> softTrim(long offsetIncluded);
LogStorage logStorage();
void cacheIndex(EntryIndex index);
void invalidateCache();
void invalidateCache(long offset);
long getMessageCount(long startOffset, long endOffset); // -1 if not available
void fence();
CompletableFuture<Void> closeAsync(); // default
CompletableFuture<Long> computeRetentionTrimOffset(long maxOffset, long retentionMillis,
long retentionSizeBytes); // default
CompletableFuture<Long> binarySearchOffset(long min, long max,
Predicate<LogEntryHeader> predicate); // defaultgetFirstOffset and getLastOffset must observe the log's current state on every call rather than serving a remembered value: a log can be written and trimmed by several holders at once, so no single handle has seen its whole history. getMessageCount uses cached index data and returns -1 when it isn't available rather than computing it another way; its Javadoc otherwise defines it as counting messages between startOffset and endOffset, and readEntries's includeTrimmed overload is documented to include soft-trimmed entries when the flag is true — see Implementation status for two documented deviations from these two contracts. What readEntry returns for an offset with no entry is a genuine gap, not a deviation: this interface leaves it undocumented; see Implementation status for how a specific implementation behaves.
Entries and offsets
One append call creates one entry, and its offset equals its first record's offset — not a batch sequence number. An entry at offset o containing n records (n always ≥ 1) spans offsets [o, o+n). LogOffset and EntryHeader/LogEntryHeader describe an entry by its offset, record count, timestamp, entry size, and cumulative size, without carrying its payload; LogOffset.NOT_FOUND is a sentinel value distinct from a valid zero offset. See the Storage Spec's Logs, entries and offsets for the normative layout this describes, and Appending for how an implementation makes an append durable and visible.
Read starts are inclusive (startOffset). Index ranges (readIndexRange, getEntryMetadataRange) and hard-trim boundaries are exclusive at the upper end. Soft trim (softTrim) marks entries up to and including an inclusive offset as deleted and returns the resulting first offset (Ursa 1.0 deviates here; see Implementation status); hard trim, in the Javadoc's words, physically removes entries up to (exclusive) a given offset, although in format version 3 it removes index records only and objects are reclaimed separately. See the Storage Spec's Reading for how a read resolves against trim markers and index records, and Trim, deletion and reclamation for the durability rules behind soft and hard trim.
Log.readEntries bounds its read by maxMessageCount — a count of records, not entries. LogCursor.readEntries (see Cursors) instead bounds its read by maxEntries, a count of entries. The two are not interchangeable units.
LogStorage.readIndexRange and Log.getEntryIndex/readIndexRange return or accept EntryIndex, and LogStorage.preFetchEntries takes a list of Position values directly as a parameter — this interface is not restricted to LogId/LogEntry/LogOffset alone. The readEntriesByIndex Javadoc says only RAW entries should be passed.
Fencing
Log.fence()'s Javadoc says only "Fences this log — subsequent append operations will fail." The interface does not say whether fencing applies to the calling handle alone, to every handle on the same log ID, or something else, and it does not say whether the fenced state is durable or held only in process memory. LogState's two values, NORMAL and FENCED, are tracked per log ID by LogStateManager, whose methods are individually thread-safe but do not make a check of the current state followed by an action on it atomic. There is no un-fence operation on this interface.
This Log.fence() mechanism is separate from the Storage Spec's durable write fences, which a catalog installs when it deletes a log — see Write leases and fences. Neither one is a distributed leader-election or ownership protocol; see Implementation status for what implementing this API does and does not establish.
Cancellation
The interfaces do not specify what happens to storage state if a caller cancels a CompletableFuture returned by append, readEntries, or another I/O method before it completes. Treat cancellation semantics as unspecified unless a specific implementation documents them.
Defaults
| Method | Default behavior |
|---|---|
LogStorage.readEntriesByIndex(logId, indices, startOffset, maxOffset, maxMessageCount, maxSizeBytes, offsetDeleted, skipCondition) | Returns an empty list if indices is empty; otherwise ignores indices, maxOffset, and both predicates, and calls readEntries(logId, startOffset, maxMessageCount, maxSizeBytes) |
LogStorage.getFirstOffset(logId, includeTrimmed) | Ignores includeTrimmed and delegates to getFirstOffset(logId) |
LogStorage.preFetchEntries(logId, positions) | No-op |
Log.closeAsync() | Calls close() synchronously on the calling thread and returns an already-completed or already-failed future; never throws |
Log.computeRetentionTrimOffset(maxOffset, retentionMillis, retentionSizeBytes) | Returns maxOffset itself, ignoring both retention arguments |
Log.binarySearchOffset(min, max, predicate) | Returns min, without searching |
computeRetentionTrimOffset and binarySearchOffset are two of the six defaults across this API that return a plausible-looking answer instead of failing — the other four are on Cursors. computeRetentionTrimOffset's default returns maxOffset itself, the answer when both retention arguments are 0 (none), not the answer for "no retention trim is needed." A caller that trims a log to this result removes every entry up to maxOffset, even when the retention arguments asked for infinite retention. binarySearchOffset's default returns min without searching, which reads as "found at the low end of the range" rather than "not implemented."
Buffer ownership
The contract uses Netty ByteBuf, not ByteBuffer:
- Append: the caller owns its reference until the returned future completes, on success or failure, then must release it exactly once. Do not release or mutate the payload while the operation is using it.
- Read: close every returned
LogEntry. Itspayload()is a borrowed, read-only view owned by that entry; do not release the view directly. - Transfer: to retain a payload beyond the entry's lifetime, take a
retainedDuplicate()of it and release that independently owned reference later.
// log is an already-open Log; this blocking example is not for an event loop.
ByteBuf data = Unpooled.copiedBuffer("hello", StandardCharsets.UTF_8);
LogEntryHeader written;
try {
written = log.append(1, data).join();
} finally {
data.release();
}
try (LogEntry entry = log.readEntry(written.offset()).join()) {
System.out.println(entry.payload().toString(StandardCharsets.UTF_8));
}The example uses io.netty.buffer.ByteBuf, io.netty.buffer.Unpooled, java.nio.charset.StandardCharsets, and io.lakestream.api types. Implementations must make repeated close() calls on the same entry safe, but callers should still structure ownership around exactly one close per returned entry.
Types
| Type | Kind | Shape |
|---|---|---|
LogId | record | (long id); wraps the numeric log identifier; of(id) is a factory |
LogEntry | interface, extends AutoCloseable | offset(), numberOfRecords() (≥ 1), timestamp(), size(), payload() returning ByteBuf, and close() |
LogEntryHeader | interface | offset(), numberOfRecords(), timestamp(), entrySize(), cumulativeSize(), newerThan(other) |
EntryHeader | record, implements LogEntryHeader | (long offset, int numberOfMessages, long writtenTimestamp, int entrySize, long cumulativeSize); NOT_FOUND is an all-zero/-1 sentinel; bridges numberOfRecords()/timestamp()/newerThan() onto its own fields |
LogEntryIndex | interface | header(), entryCount(), searchEntryHeader(offset), getFirstEntryHeader(), getLastEntryHeader() |
EntryIndex | class, implements LogEntryIndex | Holds header, position, entryCount, and indexType as private final fields, plus extraData, also private final. entryOffsets and entryHeaders are the exception: both are mutable, volatile, and — unlike the records elsewhere in this API — exposed through public setters. NOT_FOUND is a sentinel; of(header, position, entryCount, storageObjectCount) picks IndexType.COMPACT when there is exactly one storage object; getEntryHeader(long offset) looks up a specific offset's header |
EntryIndex.HeaderWithIndex | nested record | (EntryHeader header, int index) — a header paired with its position in the offsets array |
EntryIndex.IndexType | nested enum | NORMAL, COMPACT |
LogOffset | record | (long offset, int numberOfRecords, long timestamp, int entrySize, long cumulativeSize); a 3-argument constructor leaves both size fields at 0; NOT_FOUND is a sentinel |
Position | record | (FileInfo file, long indexId, FileType fileType) — an object-storage location paired with an index ID and file type; see the Storage Spec's Compacted objects for where a PARQUET Position points. Three more constructors exist: Position(String location) (indexId -1, FileType.RAW), Position(String location, long indexId, FileType fileType) (file size 0), and Position(String location, long fileSize, long indexId, FileType fileType). parseV1Format/parseV2Format parse the "location-indexId" and "location-indexId-fileType" string forms back into a Position; toV1Format()/toV2Format() render the same two forms, and toString(formatVersion) dispatches to one of them, accepting only 1 or 2 and throwing IllegalArgumentException otherwise; withoutIndexId(), location(), size(), isBinary() are also defined; NOT_FOUND is a sentinel |
Position.FileType | nested enum | RAW, PARQUET |
FileInfo | record | (String location, long size) |
LogState | enum | NORMAL (any operation allowed), FENCED (appends fail) |
LogStateManager | interface | setState(long streamId, LogState state), getState(long streamId); keyed by the raw long log/stream ID rather than LogId |
EntryHeader and EntryIndex are this API's own implementations of LogEntryHeader and LogEntryIndex — both types live in lakestream-api, in the same package as the interfaces they implement.