Skip to content

Instantly share code, notes, and snippets.

@igalshilman
Created September 3, 2026 16:18
Show Gist options
  • Select an option

  • Save igalshilman/391ca719f3834cad4716b86da87e92c4 to your computer and use it in GitHub Desktop.

Select an option

Save igalshilman/391ca719f3834cad4716b86da87e92c4 to your computer and use it in GitHub Desktop.
An example of how to asynchronously offload a long conversation to a cold storage.
/**
* Standalone example: lazy conversation history with cold-prefix offload.
*
* ## What is stored
*
* Metadata and message segments are stored as independent values. `append`
* and `offloadPrefix` opt into lazy state because they touch only selected
* segments; `applyOffload` and `history` use eager state:
*
* - `history/meta` holds `messageSequence`, `headSegment`, `tailSegment`, and
* an optional `PendingOffload`. `messageSequence` is `0` for an empty
* history and otherwise identifies the latest append.
* - `history/segment/<n>` holds one `{messages: string[]}` value. Only the tail
* is mutable. Filling it advances `tailSegment`, making the full segment
* immutable and eligible for offload.
* - `.history-archive/conversation-history/<conversation>/<segment>.json` is
* one cold file containing that segment's message array. The file body needs
* no pointer, sequence, or segment metadata because its path supplies
* identity and ordering.
*
* The application owns the archive root as process-level configuration. A
* conversation's directory is derived from its Virtual Object key, so no
* archive location is stored per conversation. The segment capacity is also
* an internal storage detail and is not part of the read contract. This
* example targets roughly 1,024 recent messages—approximately 250k tokens
* when a message averages 250 tokens.
*
* ## Offload protocol
*
* ```text
* append (exclusive)
* ├─ append to the mutable tail
* ├─ advance messageSequence
* ├─ when full, advance the tail
* └─ reserve an immutable range and send offloadPrefix
* │
* ▼
* offloadPrefix (shared)
* ├─ derive the conversation archive directory from the object key
* ├─ read the reserved immutable segments
* ├─ restate.run(write one file per segment)
* └─ send applyOffload
* │
* ▼
* applyOffload (exclusive)
* ├─ advance the committed cold/hot boundary
* └─ clear the archived Restate segment keys
* ```
*
* The archive write runs in a shared handler, so exclusive appends can continue
* while it is in flight. `PendingOffload` reserves a stable, sealed range and
* prevents overlapping writes. If several segments accumulate, one
* `restate.run` writes multiple files—still one file per segment. Deterministic
* paths and atomic renames make a partially completed batch safe to retry.
*
* The local filesystem is deliberately only an inspectable stand-in for object
* storage. It is not shared by service replicas and may disappear with a
* container. Replacing the atomic file write with S3 `PutObject` does not
* change the Restate protocol around it.
*
* ## Read contract
*
* `history()` does not read archived files. It returns:
*
* - `archive: null` until the first offload commits, otherwise the derived
* `{directory, messageCount}` where `messageCount` is the number of committed
* archived messages computed as `messageSequence - recent.length`; and
* - `recent`, the flattened suffix read from `headSegment..tailSegment` in
* Restate state.
*
* To reconstruct the complete conversation, a reader lists the returned
* directory, concatenates files in filename order until it has read
* `archive.messageCount` messages, and finally appends `recent`.
*
* A settled conversation with three archived and eight recent segments looks
* like this:
*
* ```text
* archive .history-archive/conversation-history/<conversation>/
* 000000000000.json 000000000001.json 000000000002.json
* segment 0 segment 1 segment 2
*
* Restate segment 3 ... segment 10
* meta { messageSequence: 1281, head: 3, tail: 10, offload: null }
* response { archive: { messageCount: 384 }, recent: messages 385..1281 }
* ```
*
* ## Try it
*
* Save this file as `conversation-history.ts`. With a local Restate server
* already running, install its standalone dependencies and start the endpoint
* from the directory where you want `.history-archive` to be created:
*
* ```shell
* pnpm add @restatedev/restate-sdk @restatedev/restate-sdk-gen
* pnpm add --save-dev tsx
* pnpm exec tsx conversation-history.ts
* ```
*
* In a second terminal, register the endpoint and populate a fresh Virtual
* Object key:
*
* ```shell
* restate deployments register http://localhost:9080
*
* for i in {1..1281}; do
* curl -s http://localhost:8080/ConversationHistoryExample/archive-listing-demo/append \
* -H 'content-type: application/json' \
* -d "\"Berlin itinerary note $i\""
* done
*
* curl -s -X POST \
* http://localhost:8080/ConversationHistoryExample/archive-listing-demo/history
* ```
*
* Abbreviated `history()` output:
*
* ```json
* {
* "archive": {
* "directory": ".history-archive/conversation-history/YXJjaGl2ZS1saXN0aW5nLWRlbW8",
* "messageCount": 384
* },
* "recent": ["Berlin itinerary note 385", "...", "Berlin itinerary note 1281"]
* }
* ```
*
* The returned directory is relative to the endpoint's working directory.
* Inspecting it shows three immutable segment files:
*
* ```shell
* ls .history-archive/conversation-history/YXJjaGl2ZS1saXN0aW5nLWRlbW8
* # 000000000000.json
* # 000000000001.json
* # 000000000002.json
*
* jq -s 'add | {messageCount: length, first: .[0], last: .[-1]}' \
* .history-archive/conversation-history/YXJjaGl2ZS1saXN0aW5nLWRlbW8/*.json
* # {
* # "messageCount": 384,
* # "first": "Berlin itinerary note 1",
* # "last": "Berlin itinerary note 384"
* # }
* ```
*/
import {mkdir, rename, writeFile} from "node:fs/promises";
import {join} from "node:path";
import * as restate from "@restatedev/restate-sdk-gen";
import {serve} from "@restatedev/restate-sdk";
const MESSAGES_PER_SEGMENT = 128;
const RECENT_SEGMENTS_TO_KEEP = 8;
const META_KEY = "history/meta";
// Application-owned archive configuration; no location is stored in VO state.
const ARCHIVE_ROOT = ".history-archive";
const ARCHIVE_NAMESPACE = "conversation-history";
// ---------------------------------------------------------------------------
// Explicit wire and state types. The generator SDK's default JSON serde is
// sufficient for this example, so these do not need runtime schemas.
// ---------------------------------------------------------------------------
/** One independently loaded conversation segment in Restate lazy state. */
export type HistorySegment = {
messages: string[];
};
/** Durable reservation for a range being copied to the cold archive. */
export type PendingOffload = {
firstSegment: number;
lastSegment: number;
};
/** Small routing record; archive location is derived rather than persisted. */
export type HistoryMeta = {
/** Sequence assigned to the latest append, or zero before the first one. */
messageSequence: number;
headSegment: number;
tailSegment: number;
offload: PendingOffload | null;
};
/** Derived archive location, committed message count, and recent suffix. */
export type HistoryView = {
archive: {
directory: string;
messageCount: number;
} | null;
recent: string[];
};
/** Lazy, segmented conversation-history Virtual Object. */
export const ConversationHistory = restate.object({
name: "ConversationHistoryExample",
handlers: {
/** Appends one message and starts an offload when the suffix is large. */
*append(message: string): restate.Operation<void> {
const meta = yield* readMeta();
const tail = yield* readSegment(meta.tailSegment);
tail.messages.push(message);
meta.messageSequence += 1;
// Every append changes the tail and its conversation-wide sequence.
// Segment routing changes only on rollover.
writeSegment(meta.tailSegment, tail);
if (tail.messages.length < MESSAGES_PER_SEGMENT) {
restate.state().set(META_KEY, meta);
return;
}
meta.tailSegment += 1;
if (!needsOffload(meta)) {
restate.state().set(META_KEY, meta);
return;
}
const id = conversationKey();
const offload = planOffload(meta);
// Persist the reservation before the shared offload handler can observe
// it.
meta.offload = offload;
restate.state().set(META_KEY, meta);
restate
.sendClient(ConversationHistory, id)
.offloadPrefix(offload);
},
/**
* Reads a sealed range and writes one deterministic file per segment.
*
* This shared handler derives the same archive directory as `history()` and
* cannot mutate VO state. After the external write, it durably sends the
* reserved range to the exclusive commit handler.
*/
*offloadPrefix(offload: PendingOffload): restate.Operation<void> {
const id = conversationKey();
const meta = yield* readMeta();
// A stale or superseded invocation must not write another range.
if (!sameOffload(meta.offload, offload)) {
return;
}
const segments: string[][] = [];
// The reserved range excludes the mutable tail, so these values are
// stable.
for (
let index = offload.firstSegment;
index <= offload.lastSegment;
index += 1
) {
const segment = yield* readSegment(index);
segments.push(segment.messages);
}
yield* restate.run(
({signal}) =>
putArchivedSegments(
archiveDirectory(id),
offload.firstSegment,
segments,
signal,
),
{name: "put-history-segments"},
);
yield* restate
.sendClient(ConversationHistory, id)
.applyOffload(offload);
},
/** Commits an archived range and reclaims its Restate segment keys. */
*applyOffload(offload: PendingOffload): restate.Operation<void> {
const meta = yield* readMeta();
// Only the completion for the current durable reservation may reclaim
// data.
if (!sameOffload(meta.offload, offload)) {
return;
}
// This exclusive handler commits the new cold/hot boundary and Restate
// state reclamation together.
for (
let index = offload.firstSegment;
index <= offload.lastSegment;
index += 1
) {
restate.state().clear(segmentKey(index));
}
meta.headSegment = offload.lastSegment + 1;
meta.offload = null;
restate.state().set(META_KEY, meta);
},
/** Returns committed archive metadata and the recent Restate-held suffix. */
*history(): restate.Operation<HistoryView> {
const id = conversationKey();
const meta = yield* readMeta();
const recent: string[] = [];
for (
let index = meta.headSegment;
index <= meta.tailSegment;
index += 1
) {
recent.push(...(yield* readSegment(index)).messages);
}
// Archived files can become visible before applyOffload runs. The
// exclusive commit advances headSegment only after every file in the
// range is complete. Subtracting the recent Restate suffix from the
// latest sequence therefore yields the committed archive size.
const archivedMessageCount = meta.messageSequence - recent.length;
const archive =
archivedMessageCount > 0
? {
directory: archiveDirectory(id),
messageCount: archivedMessageCount,
}
: null;
return {archive, recent};
},
},
options: {
handlers: {
append: {enableLazyState: true},
offloadPrefix: {shared: true, enableLazyState: true},
history: {shared: true},
},
},
});
// This example is its own deployable endpoint rather than part of the agent
// runtime endpoint, so it can be launched and explored independently.
serve({services: [ConversationHistory]});
// ---------------------------------------------------------------------------
// Lazy state mechanics
// ---------------------------------------------------------------------------
function* readMeta(): restate.Operation<HistoryMeta> {
// Lazy state has no physical metadata value before the first append.
return (
(yield* restate.sharedState().get<HistoryMeta>(META_KEY)) ?? {
messageSequence: 0,
headSegment: 0,
tailSegment: 0,
offload: null,
}
);
}
function* readSegment(index: number): restate.Operation<HistorySegment> {
// A newly advanced tail is represented by an absent key until its first append.
return (
(yield* restate
.sharedState()
.get<HistorySegment>(segmentKey(index))) ?? {messages: []}
);
}
function writeSegment(index: number, segment: HistorySegment): void {
restate.state().set(segmentKey(index), segment);
}
function segmentKey(index: number): string {
return `history/segment/${index}`;
}
/** Whether another immutable range must be moved out of Restate state. */
function needsOffload(meta: HistoryMeta): boolean {
const recentSegments = meta.tailSegment - meta.headSegment + 1;
return !meta.offload && recentSegments > RECENT_SEGMENTS_TO_KEEP;
}
/** Plans an offload of all excess sealed segments, but never the mutable tail. */
function planOffload(meta: HistoryMeta): PendingOffload {
const lastSegment = meta.tailSegment - RECENT_SEGMENTS_TO_KEEP;
return {
firstSegment: meta.headSegment,
lastSegment,
};
}
function sameOffload(
current: PendingOffload | null,
candidate: PendingOffload,
): boolean {
return (
current?.firstSegment === candidate.firstSegment &&
current.lastSegment === candidate.lastSegment
);
}
// ---------------------------------------------------------------------------
// Local archive boundary
// ---------------------------------------------------------------------------
async function putArchivedSegments(
directory: string,
firstSegment: number,
segments: string[][],
signal: AbortSignal,
): Promise<void> {
signal.throwIfAborted();
await mkdir(directory, {recursive: true});
for (const [offset, messages] of segments.entries()) {
signal.throwIfAborted();
const segment = firstSegment + offset;
const file = join(
directory,
`${segment.toString().padStart(12, "0")}.json`,
);
const temporaryFile = `${file}.tmp`;
const blob = JSON.stringify(messages);
await writeFile(temporaryFile, blob, {encoding: "utf8", signal});
signal.throwIfAborted();
await rename(temporaryFile, file);
console.log("archived history segment", file);
}
}
/** Reconstructs the stable archive directory from config and the VO key. */
function archiveDirectory(conversationId: string): string {
const id = Buffer.from(conversationId).toString("base64url");
return join(ARCHIVE_ROOT, ARCHIVE_NAMESPACE, id);
}
function conversationKey(): string {
const key = restate.handlerRequest().key;
return key!;
}
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment