返回源码地图

packages/experimental/browser-use-runtime/src/index.ts

main snapshot · da00f7f5358f · 正文引用章节 12;完整原文可核对,不声称全文件人工逐行审计

完整原文供逐行核对;页面收录不代表每行都经过人工语义审核。MIT 许可见 许可证。

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