1
/**2
* Browser resource ownership for the experimental providers. Resources belong3
* to an exact live Agent activation and never transfer to a resumed Session.4
* @module5
*/7
import type { Context } from '@deepseek-ai/cordis'8
import type { Agent } from '@deepseek-ai/dsh-agent'10
/** One provider-owned browser or connection and its quiescent cleanup. */11
export interface OwnedSessionResource<T> {12
/** Provider-private handle exposed to operations. */13
value: T14
/** Stop admission, interrupt pending operations, and await resource shutdown. */15
close: () => Promise<void>16
}18
/** Resource creation and attachment ownership selected by one provider. */19
export interface SessionResourceOptions<T> {20
/** Provider name included in lifecycle diagnostics. */21
label: string22
/** Reserve one existing browser for at most one live Session. */23
exclusive: boolean24
/**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
}33
interface Entry<T> {34
controller: AbortController35
ready: Promise<OwnedSessionResource<T>>36
tail: Promise<void>37
closing?: Promise<void>38
}40
/** Stop a caller's wait while retaining handlers on the resource owner's work. */41
function 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
}55
/**56
* Lazily acquires one resource per live Session and serializes its operations.57
* Provider disposal closes connections before awaiting operations, allowing58
* transport closure to interrupt work whose upstream API has no abort support.59
*/60
export 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> | undefined66
/**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>) {}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) === agent80
&& (this.entries.has(agent) || !this.options.exclusive || this.entries.size === 0)81
}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.value96
}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 } | undefined113
if (reason?.kind !== 'disposed') return114
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 result128
}).finally(() => { signal.removeEventListener('abort', releaseDisposed) })129
// The queue tracks settlement independently of a caller observing its error.130
entry.tail = task.then(() => {}, () => {})131
return task132
}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
}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 current154
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 error176
}),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 entry183
}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.tail193
}194
this.entries.delete(agent)195
})196
}197
}