Storage
SQLite
Provide PersistentHistory.layer with a SQLite store to retain completed runs across Node.js restarts:
import { import SqliteThreadStore
SqliteThreadStore } from "@yielded/agent-storage-sqlite";import { import AgentRuntime
AgentRuntime, import PersistentHistory
PersistentHistory } from "@yielded/agent";import { import Effect
Effect, import Layer
Layer } from "effect";
const const History: Layer.Layer<ThreadHistory, SqliteThreadStore.SqliteStorageInitializationError, never>
History = import PersistentHistory
PersistentHistory.const layer: Layer.Layer<ThreadHistory, never, ThreadStore>
Provide retained history to the normal AgentRuntime entry points. Each successful Run appends
one three-record batch. Staging is private to that Run; interruption discards it. Epoch zero
and the loaded tail fence stale writers without replaying external execution. A storage failure
after append may leave the whole Run recorded. Adapter failpoints cover both durable mutations.
layer.Pipeable.pipe<Layer.Layer<ThreadHistory, never, ThreadStore>, Layer.Layer<ThreadHistory, SqliteThreadStore.SqliteStorageInitializationError, never>>(this: Layer.Layer<ThreadHistory, never, ThreadStore>, ab: (_: Layer.Layer<ThreadHistory, never, ThreadStore>) => Layer.Layer<ThreadHistory, SqliteThreadStore.SqliteStorageInitializationError, never>): Layer.Layer<...> (+21 overloads)
pipe( import Layer
Layer.const provide: <never, SqliteThreadStore.SqliteStorageInitializationError, ThreadStore | ThreadReader | ThreadImport>(that: Layer.Layer<ThreadStore | ThreadReader | ThreadImport, SqliteThreadStore.SqliteStorageInitializationError, never>) => <RIn2, E2, ROut2>(self: Layer.Layer<ROut2, E2, RIn2>) => Layer.Layer<ROut2, SqliteThreadStore.SqliteStorageInitializationError | E2, Exclude<...>> (+3 overloads)
Feeds the output services of the dependency layer into the requirements of
this layer, returning a layer that only provides the services from this layer.
When to use
Use when you need to hide an implementation dependency layer from callers.
Details
In serviceLayer.pipe(Layer.provide(dependencyLayer)), the dependency layer is
built first and is used to satisfy the requirements of serviceLayer.
Example (Providing layer dependencies)
import { Context, Effect, Layer } from "effect"
class Database extends Context.Service<Database, { readonly query: (sql: string) => Effect.Effect<string>}>()("Database") {}
class UserService extends Context.Service<UserService, { readonly getUser: (id: string) => Effect.Effect<{ id: string name: string }>}>()("UserService") {}
class Logger extends Context.Service<Logger, { readonly log: (msg: string) => Effect.Effect<void>}>()("Logger") {}
// Create dependency layersconst databaseLayer = Layer.succeed(Database, { query: Effect.fn("Database.query")((sql: string) => Effect.succeed(`DB: ${sql}`))})
const logs: Array<string> = []const loggerLayer = Layer.succeed(Logger, { log: Effect.fn("Logger.log")((msg: string) => Effect.sync(() => logs.push(`[LOG] ${msg}`)))})
// UserService depends on Database and Loggerconst userServiceLayer = Layer.effect(UserService, Effect.gen(function*() { const database = yield* Database const logger = yield* Logger
return { getUser: Effect.fn("UserService.getUser")(function*(id: string) { yield* logger.log(`Looking up user ${id}`) const result = yield* database.query( `SELECT * FROM users WHERE id = ${id}` ) return { id, name: result } }) }}))
// Provide dependencies to UserService layerconst userServiceWithDependencies = userServiceLayer.pipe( Layer.provide(Layer.mergeAll(databaseLayer, loggerLayer)))
// Now UserService layer has no dependenciesconst program = Effect.gen(function*() { const userService = yield* UserService return yield* userService.getUser("123")}).pipe( Effect.provide(userServiceWithDependencies))Effect.runSync(program) // => { id: "123", name: "DB: SELECT * FROM users WHERE id = 123" }logs // => ["[LOG] Looking up user 123"]
provide(import SqliteThreadStore
SqliteThreadStore.const layer: (options: SqliteThreadStore.SqliteStorageOptions) => Layer.Layer<ThreadStore | ThreadReader | ThreadImport, SqliteThreadStore.SqliteStorageInitializationError>
A composition-root convenience Layer for canonical Threads. Durable accepted work is
served by the separate SubmissionLedger port.
layer({ SqliteStorageOptions.filename: string
filename: "./history.sqlite" })),);
const const conversation: Effect.Effect<{ readonly threadId: string & Brand<"@effect-agent/core/ThreadId">; readonly runId: string & Brand<"@effect-agent/core/RunId">; readonly output: { readonly itinerary: readonly string[]; }; readonly finishReason: "completed" | "budget-exhausted" | "model-stop"; readonly turns: number; readonly usage?: RunTotals | undefined; readonly runDisposition?: Json | undefined; readonly exhausted?: "tokens" | "tool-calls" | "turns" | undefined; readonly delegatedUsage?: RunTotals | undefined;}, SqliteThreadStore.SqliteStorageInitializationError | AgentRuntime.AgentRuntimeFailure<...>, ModelServices>
conversation = import Effect
Effect.const gen: <Effect.Effect<{ readonly threadId: string & Brand<"@effect-agent/core/ThreadId">; readonly runId: string & Brand<"@effect-agent/core/RunId">; readonly output: { readonly itinerary: readonly string[]; }; readonly finishReason: "completed" | "budget-exhausted" | "model-stop"; readonly turns: number; readonly usage?: RunTotals | undefined; readonly runDisposition?: Json | undefined; readonly exhausted?: "tokens" | "tool-calls" | "turns" | undefined; readonly delegatedUsage?: RunTotals | undefined;}, AgentRuntime.AgentRuntimeFailure<Definition<String, ... 5 more ..., undefined> & { ...;}, never, never>, AgentRuntime.AgentRuntimeRequirements<...>>, { readonly threadId: string & Brand<"@effect-agent/core/ThreadId">; readonly runId: string & Brand<"@effect-agent/core/RunId">; readonly output: { readonly itinerary: readonly string[]; }; readonly finishReason: "completed" | "budget-exhausted" | "model-stop"; readonly turns: number; readonly usage?: RunTotals | undefined; readonly runDisposition?: Json | undefined; readonly exhausted?: "tokens" | "tool-calls" | "turns" | undefined; readonly delegatedUsage?: RunTotals | undefined;}>(f: () => Generator<...>) => Effect.Effect<...> (+1 overload)
Provides a way to write effectful code using generator functions, simplifying
control flow and error handling.
When to use
Use when you want to write effectful code that looks and behaves like
synchronous code, while still handling asynchronous tasks, errors, and complex
control flow such as loops and conditions.
Generator functions work similarly to async/await but keep errors,
requirements, and interruption in the Effect type. You can yield* values
from effects and return the final result at the end.
Example (Sequencing effects with generators)
import { Data, Effect } from "effect"
class DiscountRateError extends Data.TaggedError("DiscountRateError")<{}> {}
const addServiceCharge = (amount: number) => amount + 1
const applyDiscount = ( total: number, discountRate: number): Effect.Effect<number, DiscountRateError> => discountRate === 0 ? Effect.fail(new DiscountRateError()) : Effect.succeed(total - (total * discountRate) / 100)
const fetchTransactionAmount = Effect.promise(() => Promise.resolve(100))
const fetchDiscountRate = Effect.promise(() => Promise.resolve(5))
export const program = Effect.gen(function*() { const transactionAmount = yield* fetchTransactionAmount const discountRate = yield* fetchDiscountRate const discountedAmount = yield* applyDiscount( transactionAmount, discountRate ) const finalAmount = addServiceCharge(discountedAmount) return `Final amount to charge: ${finalAmount}`})
await Effect.runPromise(program) // => "Final amount to charge: 96"
gen(function* () { const const first: { readonly threadId: string & Brand<"@effect-agent/core/ThreadId">; readonly runId: string & Brand<"@effect-agent/core/RunId">; readonly output: { readonly itinerary: readonly string[]; }; readonly finishReason: "completed" | "budget-exhausted" | "model-stop"; readonly turns: number; readonly usage?: RunTotals | undefined; readonly runDisposition?: Json | undefined; readonly exhausted?: "tokens" | "tool-calls" | "turns" | undefined; readonly delegatedUsage?: RunTotals | undefined;}
first = yield* import AgentRuntime
AgentRuntime.run<Definition<String, Struct<{ readonly itinerary: $Array<String>;}>, "Plan a trip with one itinerary entry per day.", Toolkit<{}>, undefined, undefined, undefined> & { readonly id: Brand<"@effect-agent/core/AgentId"> & "trip-planner";}, never, never>(agent: Definition<String, Struct<{ readonly itinerary: $Array<String>;}>, "Plan a trip with one itinerary entry per day.", Toolkit<{}>, undefined, undefined, undefined> & { readonly id: Brand<"@effect-agent/core/AgentId"> & "trip-planner";}, input: string, options?: RunOptions<...> | undefined): Effect.Effect<...>export run
Accept schema-encoded input, retaining runtime validation. Use runUnknown for external data.
run(const planner: Definition<String, Struct<{ readonly itinerary: $Array<String>;}>, "Plan a trip with one itinerary entry per day.", Toolkit<{}>, undefined, undefined, undefined> & { readonly id: Brand<"@effect-agent/core/AgentId"> & "trip-planner";}
planner, "Plan a trip to Lisbon"); return yield* import AgentRuntime
AgentRuntime.run<Definition<String, Struct<{ readonly itinerary: $Array<String>;}>, "Plan a trip with one itinerary entry per day.", Toolkit<{}>, undefined, undefined, undefined> & { readonly id: Brand<"@effect-agent/core/AgentId"> & "trip-planner";}, never, never>(agent: Definition<String, Struct<{ readonly itinerary: $Array<String>;}>, "Plan a trip with one itinerary entry per day.", Toolkit<{}>, undefined, undefined, undefined> & { readonly id: Brand<"@effect-agent/core/AgentId"> & "trip-planner";}, input: string, options?: RunOptions<...> | undefined): Effect.Effect<...>export run
Accept schema-encoded input, retaining runtime validation. Use runUnknown for external data.
run(const planner: Definition<String, Struct<{ readonly itinerary: $Array<String>;}>, "Plan a trip with one itinerary entry per day.", Toolkit<{}>, undefined, undefined, undefined> & { readonly id: Brand<"@effect-agent/core/AgentId"> & "trip-planner";}
planner, "Make it cheaper", { RunOptions<HookError = never, HookRequirements = never>.threadId?: (string & Brand<"@effect-agent/core/ThreadId">) | undefined
Reuse a Thread identity, including retained history, instead of allocating one.
threadId: const first: { readonly threadId: string & Brand<"@effect-agent/core/ThreadId">; readonly runId: string & Brand<"@effect-agent/core/RunId">; readonly output: { readonly itinerary: readonly string[]; }; readonly finishReason: "completed" | "budget-exhausted" | "model-stop"; readonly turns: number; readonly usage?: RunTotals | undefined; readonly runDisposition?: Json | undefined; readonly exhausted?: "tokens" | "tool-calls" | "turns" | undefined; readonly delegatedUsage?: RunTotals | undefined;}
first.threadId: string & Brand<"@effect-agent/core/ThreadId">
threadId, });}).Pipeable.pipe<Effect.Effect<{ readonly threadId: string & Brand<"@effect-agent/core/ThreadId">; readonly runId: string & Brand<"@effect-agent/core/RunId">; readonly output: { readonly itinerary: readonly string[]; }; readonly finishReason: "completed" | "budget-exhausted" | "model-stop"; readonly turns: number; readonly usage?: RunTotals | undefined; readonly runDisposition?: Json | undefined; readonly exhausted?: "tokens" | "tool-calls" | "turns" | undefined; readonly delegatedUsage?: RunTotals | undefined;}, AgentRuntime.AgentRuntimeFailure<Definition<String, ... 5 more ..., undefined> & { ...;}, never, never>, AgentRuntime.AgentRuntimeRequirements<...>>, Effect.Effect<...>>(this: Effect.Effect<...>, ab: (_: Effect.Effect<...>) => Effect.Effect<...>): Effect.Effect<...> (+21 overloads)
pipe(import Effect
Effect.const provide: <ThreadHistory, SqliteThreadStore.SqliteStorageInitializationError, never>(layer: Layer.Layer<ThreadHistory, SqliteThreadStore.SqliteStorageInitializationError, never>, options?: { readonly local?: boolean | undefined;} | undefined) => <A, E, R>(self: Effect.Effect<A, E, R>) => Effect.Effect<A, SqliteThreadStore.SqliteStorageInitializationError | E, Exclude<R, ThreadHistory>> (+5 overloads)
Provides dependencies to an effect using layers or a context. Use options.local
to build the layer every time; by default, layers are shared between provide
calls.
Example (Providing dependencies with a layer)
import { Context, Effect, Layer } from "effect"
interface Database { readonly query: (sql: string) => Effect.Effect<string>}
const Database = Context.Service<Database>("Database")
const DatabaseLayer = Layer.succeed(Database)({ query: Effect.fn("Database.query")((sql: string) => Effect.succeed(`Result for: ${sql}`))})
const program = Effect.gen(function*() { const db = yield* Database return yield* db.query("SELECT * FROM users")})
const provided = Effect.provide(program, DatabaseLayer)
await Effect.runPromise(provided) // => "Result for: SELECT * FROM users"
provide(const History: Layer.Layer<ThreadHistory, SqliteThreadStore.SqliteStorageInitializationError, never>
History));Add the adapter to your existing Yielded Agent application:
bun add @yielded/agent-storage-sqlite@beta effectHere, planner is an agent definition. Supply its model and tool
services around the program. The storage Layer opens a scoped Node SQLite client
and supplies Crypto. Keep framework packages at one release and use compatible
Effect packages.
Keep the database file on persistent storage. Save the returned threadId and pass
it to later runs to continue the conversation, including after a process restart.
Provide the history Layer around the complete program or application runtime.
SqliteThreadStore.layer({ filename, synchronous: "NORMAL" }) opts this store’s connection
into WAL NORMAL synchronization; the default is FULL. Choose the mode when constructing
the store, and use the same choice for storage Layers sharing one SQL client.
WAL NORMAL survives process crashes, but power loss or an OS crash can lose acknowledged
commits and cause external effects to repeat during durable recovery.
What is retained
Section titled “What is retained”PersistentHistory.layer commits each successful run’s input and native messages
as one atomic batch before publishing RunCompleted. A failure or interruption
before commit leaves none of that run in history. A storage error after commit can
leave the whole run recorded, so inspect history before retrying.
The database retains conversation history across restarts. An interrupted run has no automatic recovery. See retained history for the commit policy, concurrency, and history inspection.
Recover unfinished work
Section titled “Recover unfinished work”Use @yielded/agent-platform-node when accepted work must
survive a restart. It assembles SQLite storage, agent registrations, admission,
recovery, and a worker pool. Run one live host per SQLite file.
For custom durable assemblies, SqliteSubmissionLedger.ledgerLayer(options)
provides the separate accepted-work ledger. Point it at the same database file as
SqliteThreadStore.layer(options) so ownership claims fence the same thread log.
SQLite tracks layout separately from record meaning. This unreleased protocol opens only fresh
layout-21 stores with effect-agent/thread@3 records. Predecessor or ambiguous stores are rejected
before mutation. Same-format export/import installs into an empty destination Thread; no older
layout upgrade or archive converter is provided. See the
operator procedure.