index.ts 8.5 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197
  1. /**
  2. * Browser resource ownership for the experimental providers. Resources belong
  3. * to an exact live Agent activation and never transfer to a resumed Session.
  4. * @module
  5. */
  6. import type { Context } from '@deepseek-ai/cordis'
  7. import type { Agent } from '@deepseek-ai/dsh-agent'
  8. /** One provider-owned browser or connection and its quiescent cleanup. */
  9. export interface OwnedSessionResource<T> {
  10. /** Provider-private handle exposed to operations. */
  11. value: T
  12. /** Stop admission, interrupt pending operations, and await resource shutdown. */
  13. close: () => Promise<void>
  14. }
  15. /** Resource creation and attachment ownership selected by one provider. */
  16. export interface SessionResourceOptions<T> {
  17. /** Provider name included in lifecycle diagnostics. */
  18. label: string
  19. /** Reserve one existing browser for at most one live Session. */
  20. exclusive: boolean
  21. /**
  22. * Acquire one resource; reject only after rolling back partial acquisition.
  23. * @param agent - exact live owner of this acquisition.
  24. * @param signal - aborts when that owner or the provider is disposed.
  25. * @returns the acquired resource and its cleanup.
  26. */
  27. open: (agent: Agent, signal: AbortSignal) => Promise<OwnedSessionResource<T>>
  28. }
  29. interface Entry<T> {
  30. controller: AbortController
  31. ready: Promise<OwnedSessionResource<T>>
  32. tail: Promise<void>
  33. closing?: Promise<void>
  34. }
  35. /** Stop a caller's wait while retaining handlers on the resource owner's work. */
  36. function awaitOperation<T>(operation: Promise<T>, signal: AbortSignal): Promise<T> {
  37. return new Promise<T>((resolve, reject) => {
  38. const aborted = (): void => { reject(signal.reason instanceof Error ? signal.reason : new Error('browser operation canceled', { cause: signal.reason })) }
  39. signal.addEventListener('abort', aborted, { once: true })
  40. void operation.then((value) => {
  41. signal.removeEventListener('abort', aborted)
  42. resolve(value)
  43. }, (error: unknown) => {
  44. signal.removeEventListener('abort', aborted)
  45. reject(error instanceof Error ? error : new Error(String(error), { cause: error }))
  46. })
  47. })
  48. }
  49. /**
  50. * Lazily acquires one resource per live Session and serializes its operations.
  51. * Provider disposal closes connections before awaiting operations, allowing
  52. * transport closure to interrupt work whose upstream API has no abort support.
  53. */
  54. export class SessionResources<T> {
  55. private readonly entries = new Map<Agent, Entry<T>>()
  56. private readonly ownerCleanups = new Map<Agent, () => Promise<void>>()
  57. private readonly disposedOwners = new WeakSet<Agent>()
  58. private disposing: Promise<void> | undefined
  59. /**
  60. * @param ctx - provider context with the live Agent registry.
  61. * @param options - provider-owned acquisition and attachment policy.
  62. */
  63. constructor(private readonly ctx: Context, private readonly options: SessionResourceOptions<T>) {}
  64. /**
  65. * Check admission without reserving or acquiring a browser.
  66. * @param agent - exact live Agent that would own the resource.
  67. * @returns whether this owner can use or acquire the configured browser.
  68. */
  69. available(agent: Agent): boolean {
  70. return this.disposing === undefined && !this.disposedOwners.has(agent)
  71. && this.ctx.get('agents')?.get(agent.id) === agent
  72. && (this.entries.has(agent) || !this.options.exclusive || this.entries.size === 0)
  73. }
  74. /**
  75. * Obtain the current activation's resource, acquiring it once when absent.
  76. * @param agent - exact live owner, never merely a durable Session id.
  77. * @param signal - optional cancellation of this wait; acquisition remains Session-owned.
  78. * @returns the provider's resource after acquisition and ownership checks.
  79. */
  80. async get(agent: Agent, signal?: AbortSignal): Promise<T> {
  81. signal?.throwIfAborted()
  82. const entry = this.entry(agent)
  83. const resource = await (signal === undefined ? entry.ready : awaitOperation(entry.ready, signal))
  84. signal?.throwIfAborted()
  85. entry.controller.signal.throwIfAborted()
  86. return resource.value
  87. }
  88. /**
  89. * Run after earlier operations on this Session settle; other Sessions proceed independently.
  90. * Cancellation stops this caller's acquisition wait without canceling Session-owned initialization.
  91. * It reaches an active provider operation and prevents queued work from starting.
  92. * @param agent - exact live resource owner.
  93. * @param signal - cancellation for this operation.
  94. * @param operation - provider call, which must retain ownership until its work settles.
  95. * @returns the operation result or its acquisition, cancellation, or execution failure.
  96. */
  97. run<R>(agent: Agent, signal: AbortSignal, operation: (resource: T, signal: AbortSignal) => Promise<R>): Promise<R> {
  98. signal.throwIfAborted()
  99. const entry = this.entry(agent)
  100. const combined = AbortSignal.any([signal, entry.controller.signal])
  101. const releaseDisposed = () => {
  102. const reason = signal.reason as { kind?: unknown } | undefined
  103. if (reason?.kind !== 'disposed') return
  104. this.disposedOwners.add(agent)
  105. // AgentHandle waits for idle before disposing its scope; close interrupts the owned operation first.
  106. void this.closeEntry(agent, entry).catch((error: unknown) => {
  107. this.ctx.logger.warn(`${this.options.label}: browser cleanup during Session cancellation failed: ${String(error)}`)
  108. })
  109. }
  110. signal.addEventListener('abort', releaseDisposed, { once: true })
  111. const task = entry.tail.then(async () => {
  112. combined.throwIfAborted()
  113. const resource = await awaitOperation(entry.ready, combined)
  114. combined.throwIfAborted()
  115. const result = await operation(resource.value, combined)
  116. combined.throwIfAborted()
  117. return result
  118. }).finally(() => { signal.removeEventListener('abort', releaseDisposed) })
  119. // The queue tracks settlement independently of a caller observing its error.
  120. entry.tail = task.then(() => {}, () => {})
  121. return task
  122. }
  123. /**
  124. * Stop new acquisitions and await every acquired resource and owned operation.
  125. * A failed close retains its entry and rejects disposal, preserving exclusive ownership.
  126. * @returns the shared quiescent disposal promise.
  127. */
  128. dispose(): Promise<void> {
  129. return this.disposing ??= Promise.resolve().then(async () => {
  130. const settled = await Promise.allSettled([...this.entries].map(([agent, entry]) => this.closeEntry(agent, entry)))
  131. const errors = settled.flatMap(result => result.status === 'rejected' ? [result.reason as unknown] : [])
  132. if (errors.length > 0) throw new AggregateError(errors, `${this.options.label}: browser cleanup failed`)
  133. await Promise.all([...this.ownerCleanups.values()].map(close => close()))
  134. })
  135. }
  136. private entry(agent: Agent): Entry<T> {
  137. if (this.disposing !== undefined || this.disposedOwners.has(agent) || this.ctx.get('agents')?.get(agent.id) !== agent) {
  138. throw new Error(`${this.options.label}: Session is not a live browser owner`)
  139. }
  140. const current = this.entries.get(agent)
  141. if (current !== undefined) return current
  142. if (this.options.exclusive && this.entries.size > 0) {
  143. throw new Error(`${this.options.label}: attached browser is already reserved by another Session`)
  144. }
  145. if (!this.ownerCleanups.has(agent)) {
  146. const cleanup = agent.ctx.effect(() => async () => {
  147. this.disposedOwners.add(agent)
  148. const owned = this.entries.get(agent)
  149. if (owned !== undefined) await this.closeEntry(agent, owned)
  150. this.ownerCleanups.delete(agent)
  151. }, `${this.options.label}.session`)
  152. this.ownerCleanups.set(agent, cleanup)
  153. }
  154. const controller = new AbortController()
  155. const entry: Entry<T> = {
  156. controller,
  157. ready: Promise.resolve().then(() => {
  158. controller.signal.throwIfAborted()
  159. return this.options.open(agent, controller.signal)
  160. }).catch((error: unknown) => {
  161. // open() owns rollback; a failed acquisition has no remaining resource.
  162. this.entries.delete(agent)
  163. throw error
  164. }),
  165. tail: Promise.resolve(),
  166. }
  167. // Acquisition can outlive every canceled caller; later consumers still receive its failure.
  168. void entry.ready.catch(() => {})
  169. this.entries.set(agent, entry)
  170. return entry
  171. }
  172. private closeEntry(agent: Agent, entry: Entry<T>): Promise<void> {
  173. return entry.closing ??= Promise.resolve().then(async () => {
  174. entry.controller.abort(new Error(`${this.options.label}: Session browser is closing`))
  175. const resource = await entry.ready.catch(() => undefined)
  176. try {
  177. await resource?.close()
  178. } finally {
  179. await entry.tail
  180. }
  181. this.entries.delete(agent)
  182. })
  183. }
  184. }