import type { KVNamespaceListOptions as CloudflareKVNamespaceListOptions, KVNamespacePutOptions as CloudflareKVNamespacePutOptions, } from "@cloudflare/workers-types"; import { Context, Data, Effect, Option, Schema as S, type Layer } from "effect"; import * as Binding from "./Binding"; import type { WorkerEnvironment } from "./Environment"; const expectedKvNamespace = "KV namespace binding with get(), put(), delete(), getWithMetadata(), and list()"; /** Error raised when a KV operation fails. */ export class KvOperationError extends Data.TaggedError("KvOperationError")<{ readonly binding: string; readonly operation: string; readonly cause: unknown; }> {} /** `KVNamespace.put` options. */ export type KvPutOptions = CloudflareKVNamespacePutOptions; /** `KVNamespace.list` options with optional metadata decoding schema. */ export type KvListOptions = CloudflareKVNamespaceListOptions & { readonly metadataSchema?: S.Codec; }; /** Successful value returned by `getWithMetadata`. */ export interface KvWithMetadata { readonly value: Value; readonly metadata: Option.Option; readonly cacheStatus: Option.Option; } /** Decoded key entry returned by `list`. */ export interface KvListKey { readonly name: Key; readonly expiration: Option.Option; readonly metadata: Option.Option; } /** Decoded result returned by `list`. */ export interface KvListResult { readonly keys: ReadonlyArray>; readonly listComplete: boolean; readonly cursor: Option.Option; readonly cacheStatus: Option.Option; } /** * Typed KV binding definition. */ export interface KvDefinition { /** Binding name as configured in `wrangler.jsonc`. */ readonly binding: string; /** Codec used to encode/decode keys. */ readonly key: S.Codec; /** Codec used to encode/decode values. */ readonly value: S.Codec; } /** * Reusable typed KV resource definition without a concrete Cloudflare binding name. */ export interface Definition< Id extends string = string, Key = unknown, Value = unknown, EncodedValue = unknown, > { readonly id: Id; /** Codec used to encode/decode keys. */ readonly key: S.Codec; /** Codec used to encode/decode values. */ readonly value: S.Codec; } export namespace Definition { export type Any = Definition; } export interface KvClient { readonly put: ( key: Key, value: Value, options?: KvPutOptions, ) => Effect.Effect; readonly get: (key: Key) => Effect.Effect, KvOperationError | S.SchemaError>; readonly getWithMetadata: ( key: Key, metadataSchema: S.Codec, ) => Effect.Effect< Option.Option>, KvOperationError | S.SchemaError >; readonly list: ( options?: KvListOptions, ) => Effect.Effect, KvOperationError | S.SchemaError>; readonly remove: (key: Key) => Effect.Effect; readonly unsafeRaw: Effect.Effect; readonly definition: KvDefinition; } export type LayerOptions = { readonly binding: string; }; export interface TagClass< Self, Id extends string, Key, Value, EncodedValue, > extends Context.ServiceClass> { readonly id: Id; readonly keySchema: S.Codec; readonly valueSchema: S.Codec; readonly layer: ( options: LayerOptions, ) => Layer.Layer< Self, Binding.BindingNotFoundError | Binding.BindingValidationError, WorkerEnvironment >; } const maybeString = (value: string | null | undefined): Option.Option => value == null ? Option.none() : Option.some(value); const maybeNumber = (value: number | undefined): Option.Option => value === undefined ? Option.none() : Option.some(value); const kvError = (binding: string, operation: string, cause: unknown) => new KvOperationError({ binding, operation, cause }); const tryKvPromise = ( binding: string, operation: string, evaluate: () => Promise, ): Effect.Effect => Effect.tryPromise({ try: evaluate, catch: (cause) => kvError(binding, operation, cause), }); export const isKvNamespace = (value: unknown): value is KVNamespace => { if (typeof value !== "object" || value === null) { return false; } const resource = value as Record; return ( typeof resource.get === "function" && typeof resource.put === "function" && typeof resource.delete === "function" && typeof resource.getWithMetadata === "function" && typeof resource.list === "function" ); }; export const makeClient = ( definition: KvDefinition, ): ((kv: KVNamespace) => KvClient) => { const encodeKey = S.encodeEffect(definition.key); const decodeKey = S.decodeUnknownEffect(definition.key); const encodeValue = S.encodeEffect(S.fromJsonString(S.toCodecJson(definition.value))); const decodeValue = S.decodeUnknownEffect(S.fromJsonString(S.toCodecJson(definition.value))); return (kv) => ({ definition, put: Effect.fnUntraced(function* (key: Key, value: Value, options?: KvPutOptions) { const keyEncoded = yield* encodeKey(key); const valueEncoded = yield* encodeValue(value); yield* tryKvPromise(definition.binding, "put", () => kv.put(keyEncoded, valueEncoded, options), ); }), get: Effect.fnUntraced(function* (key: Key) { const keyEncoded = yield* encodeKey(key); const valueEncoded = yield* tryKvPromise(definition.binding, "get", () => kv.get(keyEncoded)); if (valueEncoded === null) { return Option.none(); } return yield* decodeValue(valueEncoded).pipe(Effect.map(Option.some)); }), getWithMetadata: Effect.fnUntraced(function* ( key: Key, metadataSchema: S.Codec, ) { const keyEncoded = yield* encodeKey(key); const result = yield* tryKvPromise(definition.binding, "getWithMetadata", () => kv.getWithMetadata(keyEncoded), ); if (result.value === null) { return Option.none(); } const value = yield* decodeValue(result.value); const metadata = result.metadata === null ? Option.none() : Option.some(yield* S.decodeUnknownEffect(metadataSchema)(result.metadata)); return Option.some({ value, metadata, cacheStatus: maybeString(result.cacheStatus), }); }), list: Effect.fnUntraced(function* (options?: KvListOptions) { const { metadataSchema, ...kvOptions } = options ?? {}; const result = yield* tryKvPromise(definition.binding, "list", () => kv.list(kvOptions), ); const keys: Array> = []; for (const key of result.keys) { const decodedName = yield* decodeKey(key.name); const decodedMetadata = key.metadata === undefined ? Option.none() : metadataSchema === undefined ? Option.some(key.metadata as Metadata) : Option.some(yield* S.decodeUnknownEffect(metadataSchema)(key.metadata)); keys.push({ name: decodedName, expiration: maybeNumber(key.expiration), metadata: decodedMetadata, }); } return { keys, listComplete: result.list_complete, cursor: maybeString("cursor" in result ? result.cursor : undefined), cacheStatus: maybeString(result.cacheStatus), }; }), remove: Effect.fnUntraced(function* (key: Key) { const keyEncoded = yield* encodeKey(key); yield* tryKvPromise(definition.binding, "delete", () => kv.delete(keyEncoded)); }), unsafeRaw: Effect.succeed(kv), }); }; export const layer = ( tag: Context.Service>, definition: KvDefinition, ) => Binding.layer( tag, definition.binding, isKvNamespace, makeClient(definition), { expected: expectedKvNamespace }, ); const makeDefinition = ( id: Id, definition: { readonly key: S.Codec; readonly value: S.Codec }, ) => { type SelfDefinition = Definition; const kvDefinition: SelfDefinition = { id, key: definition.key, value: definition.value, }; return Object.assign(kvDefinition, { layer: ( tag: Context.Service>, binding: LayerOptions, ) => layer(tag, { ...binding, key: definition.key, value: definition.value, }), }); }; export const Tag = () => ( id: Id, definition: { readonly key: S.Codec; readonly value: S.Codec; }, ) => { const kvDefinition = makeDefinition(id, definition); const tag = Context.Service>()(id); const makeLayer = (binding: LayerOptions) => kvDefinition.layer(tag, binding); return Object.assign(tag, { id: kvDefinition.id, keySchema: kvDefinition.key, valueSchema: kvDefinition.value, layer: makeLayer, }) as TagClass; };