import type * as cf from "@cloudflare/workers-types"; import * as Effect from "effect/Effect"; import * as Stream from "effect/Stream"; import type { RuntimeContext } from "../../RuntimeContext.ts"; // --------------------------------------------------------------------------- // SqlStorage — Effect-native wrapper around cf.SqlStorage // --------------------------------------------------------------------------- export type SqlStorageValue = cf.SqlStorageValue; export interface SqlCursor< T extends Record, > extends Stream.Stream { next(): Effect.Effect< { done?: false; value: T } | { done: true; value?: never }, never, RuntimeContext >; toArray(): Effect.Effect; one(): Effect.Effect; raw(): Stream.Stream; readonly columnNames: string[]; readonly rowsRead: Effect.Effect; readonly rowsWritten: Effect.Effect; } export interface SqlStorage { /** * The raw underlying Cloudflare SqlStorage binding. * * Use this when you need direct access for libraries that already support * Cloudflare Durable Object SQLite storage. */ readonly raw: cf.SqlStorage; exec>( query: string, ...bindings: any[] ): Effect.Effect, never, RuntimeContext>; readonly databaseSize: number; } const fromSqlCursor = >( cursor: cf.SqlStorageCursor, ): SqlCursor => { const stream = Stream.fromIterableEffect(Effect.sync(() => cursor)); return Object.assign(stream, { next: () => Effect.sync(() => cursor.next()), toArray: () => Effect.sync(() => cursor.toArray()), one: () => Effect.sync(() => cursor.one()), raw: () => Stream.fromIterableEffect(Effect.sync(() => cursor.raw())), get columnNames() { return cursor.columnNames; }, rowsRead: Effect.sync(() => cursor.rowsRead), rowsWritten: Effect.sync(() => cursor.rowsWritten), }) as SqlCursor; }; const fromSqlStorage = (sql: cf.SqlStorage): SqlStorage => ({ raw: sql, exec: >( query: string, ...bindings: any[] ): Effect.Effect> => Effect.sync(() => fromSqlCursor(sql.exec(query, ...bindings))), get databaseSize() { return sql.databaseSize; }, }); // --------------------------------------------------------------------------- // DurableObjectTransaction // --------------------------------------------------------------------------- export interface DurableObjectTransaction { get( key: string, options?: cf.DurableObjectGetOptions, ): Effect.Effect; get( keys: string[], options?: cf.DurableObjectGetOptions, ): Effect.Effect, never, RuntimeContext>; list( options?: cf.DurableObjectListOptions, ): Effect.Effect, never, RuntimeContext>; put( key: string, value: T, options?: cf.DurableObjectPutOptions, ): Effect.Effect; put( entries: Record, options?: cf.DurableObjectPutOptions, ): Effect.Effect; delete( key: string, options?: cf.DurableObjectPutOptions, ): Effect.Effect; delete( keys: string[], options?: cf.DurableObjectPutOptions, ): Effect.Effect; rollback(): Effect.Effect; getAlarm( options?: cf.DurableObjectGetAlarmOptions, ): Effect.Effect; setAlarm( scheduledTime: number | Date, options?: cf.DurableObjectSetAlarmOptions, ): Effect.Effect; deleteAlarm( options?: cf.DurableObjectSetAlarmOptions, ): Effect.Effect; } // --------------------------------------------------------------------------- // DurableObjectStorage // --------------------------------------------------------------------------- export interface DurableObjectStorage { get( key: string, options?: cf.DurableObjectGetOptions, ): Effect.Effect; get( keys: string[], options?: cf.DurableObjectGetOptions, ): Effect.Effect, never, RuntimeContext>; list( options?: cf.DurableObjectListOptions, ): Effect.Effect, never, RuntimeContext>; put( key: string, value: T, options?: cf.DurableObjectPutOptions, ): Effect.Effect; put( entries: Record, options?: cf.DurableObjectPutOptions, ): Effect.Effect; delete( key: string, options?: cf.DurableObjectPutOptions, ): Effect.Effect; delete( keys: string[], options?: cf.DurableObjectPutOptions, ): Effect.Effect; deleteAll( options?: cf.DurableObjectPutOptions, ): Effect.Effect; transaction( closure: ( txn: DurableObjectTransaction, ) => Effect.Effect, ): Effect.Effect; getAlarm( options?: cf.DurableObjectGetAlarmOptions, ): Effect.Effect; setAlarm( scheduledTime: number | Date, options?: cf.DurableObjectSetAlarmOptions, ): Effect.Effect; deleteAlarm( options?: cf.DurableObjectSetAlarmOptions, ): Effect.Effect; sync(): Effect.Effect; sql: SqlStorage; kv: cf.SyncKvStorage; getCurrentBookmark(): Effect.Effect; getBookmarkForTime( timestamp: number | Date, ): Effect.Effect; onNextSessionRestoreBookmark( bookmark: string, ): Effect.Effect; } // --------------------------------------------------------------------------- // Constructors from raw Cloudflare types // --------------------------------------------------------------------------- export const fromDurableObjectTransaction = ( txn: cf.DurableObjectTransaction, ): DurableObjectTransaction => ({ get: ((keyOrKeys: string | string[], options?: cf.DurableObjectGetOptions) => Effect.tryPromise(() => txn.get(keyOrKeys as any, options))) as any, list: (options?: cf.DurableObjectListOptions) => Effect.tryPromise(() => txn.list(options)), put: (( keyOrEntries: string | Record, valueOrOptions?: unknown, maybeOptions?: cf.DurableObjectPutOptions, ) => typeof keyOrEntries === "string" ? Effect.tryPromise(() => txn.put(keyOrEntries, valueOrOptions, maybeOptions), ) : Effect.tryPromise(() => txn.put( keyOrEntries, valueOrOptions as cf.DurableObjectPutOptions | undefined, ), )) as any, delete: (( keyOrKeys: string | string[], options?: cf.DurableObjectPutOptions, ) => Effect.tryPromise(() => txn.delete(keyOrKeys as any, options))) as any, rollback: () => Effect.sync(() => txn.rollback()), getAlarm: (options?: cf.DurableObjectGetAlarmOptions) => Effect.tryPromise(() => txn.getAlarm(options)), setAlarm: ( scheduledTime: number | Date, options?: cf.DurableObjectSetAlarmOptions, ) => Effect.tryPromise(() => txn.setAlarm(scheduledTime, options)), deleteAlarm: (options?: cf.DurableObjectSetAlarmOptions) => Effect.tryPromise(() => txn.deleteAlarm(options)), }); export const fromDurableObjectStorage = ( storage: cf.DurableObjectStorage, ): DurableObjectStorage => ({ get: ((keyOrKeys: string | string[], options?: cf.DurableObjectGetOptions) => Effect.tryPromise(() => storage.get(keyOrKeys as any, options))) as any, list: (options?: cf.DurableObjectListOptions) => Effect.tryPromise(() => storage.list(options)), put: (( keyOrEntries: string | Record, valueOrOptions?: unknown, maybeOptions?: cf.DurableObjectPutOptions, ) => typeof keyOrEntries === "string" ? Effect.tryPromise(() => storage.put(keyOrEntries, valueOrOptions, maybeOptions), ) : Effect.tryPromise(() => storage.put( keyOrEntries, valueOrOptions as cf.DurableObjectPutOptions | undefined, ), )) as any, delete: (( keyOrKeys: string | string[], options?: cf.DurableObjectPutOptions, ) => Effect.tryPromise(() => storage.delete(keyOrKeys as any, options))) as any, deleteAll: (options?: cf.DurableObjectPutOptions) => Effect.tryPromise(() => storage.deleteAll(options)), transaction: ( closure: (txn: DurableObjectTransaction) => Effect.Effect, ) => Effect.tryPromise(() => storage.transaction((txn) => Effect.runPromise(closure(fromDurableObjectTransaction(txn))), ), ), getAlarm: (options?: cf.DurableObjectGetAlarmOptions) => Effect.tryPromise(() => storage.getAlarm(options)), setAlarm: ( scheduledTime: number | Date, options?: cf.DurableObjectSetAlarmOptions, ) => Effect.tryPromise(() => storage.setAlarm(scheduledTime, options)), deleteAlarm: (options?: cf.DurableObjectSetAlarmOptions) => Effect.tryPromise(() => storage.deleteAlarm(options)), sync: () => Effect.tryPromise(() => storage.sync()), sql: fromSqlStorage(storage.sql), kv: storage.kv, getCurrentBookmark: () => Effect.tryPromise(() => storage.getCurrentBookmark()), getBookmarkForTime: (timestamp: number | Date) => Effect.tryPromise(() => storage.getBookmarkForTime(timestamp)), onNextSessionRestoreBookmark: (bookmark: string) => Effect.tryPromise(() => storage.onNextSessionRestoreBookmark(bookmark)), });