Created
September 3, 2026 16:18
-
-
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.
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
| /** | |
| * 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