( lease: StorageRootLease<K, 'write'>, kind: K, extension: E, )
| 268 | } |
| 269 | |
| 270 | async function createExecutionStoresForWrite<K extends StorageRootKind, E extends object>( |
| 271 | lease: StorageRootLease<K, 'write'>, |
| 272 | kind: K, |
| 273 | extension: E, |
| 274 | ): Promise<ExecutionStoresWriterBase<K> & E> { |
| 275 | const sessionStore = createSessionStore(lease.canonicalPath); |
| 276 | const agentRunStore = createSqliteAgentRunStore(lease.canonicalPath); |
| 277 | const interactionStore = |
| 278 | 'interactionStore' in extension |
| 279 | ? (extension.interactionStore as InteractiveInteractionStoreWriterFacade) |
| 280 | : undefined; |
| 281 | const runtimePersistence = await openRuntimeEventPersistence({ |
| 282 | workspaceRoot: lease.canonicalPath, |
| 283 | }).catch(async (error) => { |
| 284 | await sessionStore.close?.().catch(() => {}); |
| 285 | agentRunStore.close?.(); |
| 286 | if (interactionStore) closeSqliteInteractionStoreFacade(interactionStore); |
| 287 | throw error; |
| 288 | }); |
| 289 | const runtimeEventStore = runtimePersistence.runtimeEventStore; |
| 290 | let conversationOperationalStateStore: ConversationOperationalStateStore; |
| 291 | try { |
| 292 | conversationOperationalStateStore = createConversationOperationalStateStore( |
| 293 | lease.canonicalPath, |
| 294 | ); |
| 295 | } catch (error) { |
| 296 | await closeExecutionStorePersistence(sessionStore, runtimePersistence, { |
| 297 | agentRunStore, |
| 298 | interactionStore, |
| 299 | }).catch(() => {}); |
| 300 | throw error; |
| 301 | } |
| 302 | const messageReceiptStore = createSqliteMessageReceiptStore(lease.canonicalPath); |
| 303 | await Promise.all([agentRunStore.ready?.(), messageReceiptStore.ready()]).catch(async (error) => { |
| 304 | await closeExecutionStorePersistence(sessionStore, runtimePersistence, { |
| 305 | agentRunStore, |
| 306 | conversationOperationalStateStore, |
| 307 | messageReceiptStore, |
| 308 | interactionStore, |
| 309 | }).catch(() => {}); |
| 310 | throw error; |
| 311 | }); |
| 312 | const run = <T>(operation: () => Promise<T>) => |
| 313 | runWithStorageRootLease(lease, kind, 'write', operation); |
| 314 | |
| 315 | const stores: ExecutionStoresWriterBase<K> & E = { |
| 316 | ...extension, |
| 317 | kind, |
| 318 | [executionStoresWriterBrand]: kind, |
| 319 | purgeConversationOperationalState: (sessionId) => |
| 320 | run(() => conversationOperationalStateStore.purge(sessionId)), |
| 321 | sessionStore: { |
| 322 | ready: () => run(() => sessionStore.ready()), |
| 323 | create: (input, initialBoundary) => run(() => sessionStore.create(input, initialBoundary)), |
| 324 | createImportedSession: (input, messages) => |
| 325 | run(() => sessionStore.createImportedSession(input, messages)), |
| 326 | probeStableSessionCreate: (sessionId, requestFingerprint) => |
| 327 | run(() => sessionStore.probeStableSessionCreate(sessionId, requestFingerprint)), |
no test coverage detected