activity.ts 3.2 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687
  1. /** One-shot Team change waiters independent of durable state projection. */
  2. import type { TeamId, TeamWaitResult } from './types.ts'
  3. import { errorMessage, TeamError } from './error.ts'
  4. interface Waiter {
  5. readonly resolve: () => void
  6. }
  7. /** Owns current Team change waiters and releases each at most once. */
  8. export class TeamActivity {
  9. private readonly waiters = new Map<TeamId, Set<Waiter>>()
  10. private closed = false
  11. /**
  12. * Wait for one later Team-domain or member-status change.
  13. * @param id - Team whose next edge wakes the caller.
  14. * @param timeoutMs - bounded wait duration from ten seconds through one hour.
  15. * @param signal - caller cancellation for this wait only.
  16. * @returns whether the wait ended by timeout.
  17. */
  18. async wait(id: TeamId, timeoutMs: number, signal: AbortSignal): Promise<TeamWaitResult> {
  19. if (!Number.isSafeInteger(timeoutMs) || timeoutMs < 10_000 || timeoutMs > 3_600_000) {
  20. throw new TeamError('timeoutMs must be an integer from 10000 through 3600000', 'TEAM_INVALID_TIMEOUT')
  21. }
  22. signal.throwIfAborted()
  23. if (this.closed) return { timedOut: false }
  24. const changed = await new Promise<boolean>((resolve, reject) => {
  25. let waiters = this.waiters.get(id)
  26. if (waiters === undefined) {
  27. waiters = new Set()
  28. this.waiters.set(id, waiters)
  29. }
  30. let settled = false
  31. const finish = (settle: () => void): void => {
  32. /* v8 ignore next -- timeout, abort, and notification may race after one winner removes the others. */
  33. if (settled) return
  34. settled = true
  35. clearTimeout(timer)
  36. signal.removeEventListener('abort', onAbort)
  37. waiters.delete(waiter)
  38. if (waiters.size === 0) this.waiters.delete(id)
  39. settle()
  40. }
  41. const onAbort = (): void => {
  42. finish(() => {
  43. const reason: unknown = signal.reason
  44. reject(reason instanceof Error
  45. ? reason
  46. : new TeamError(`wait_agent aborted: ${errorMessage(reason)}`, 'TEAM_WAIT_ABORTED'))
  47. })
  48. }
  49. const waiter: Waiter = {
  50. resolve: () => {
  51. finish(() => { resolve(true) })
  52. },
  53. }
  54. waiters.add(waiter)
  55. const timer = setTimeout(() => { finish(() => { resolve(false) }) }, timeoutMs)
  56. signal.addEventListener('abort', onAbort, { once: true })
  57. // AbortSignal does not replay an abort that wins between the pre-check and listener registration.
  58. /* v8 ignore next -- requires an abort in the synchronous gap between the pre-check and listener registration. */
  59. if (signal.aborted) onAbort()
  60. })
  61. return { timedOut: !changed }
  62. }
  63. /**
  64. * Wake and remove every current waiter for one Team.
  65. * @param id - Team whose current waiters observe the change.
  66. */
  67. notify(id: TeamId): void {
  68. const waiters = this.waiters.get(id)
  69. if (waiters === undefined) return
  70. this.waiters.delete(id)
  71. for (const waiter of waiters) waiter.resolve()
  72. }
  73. /** Close admission and wake every current waiter during runtime disposal. */
  74. close(): void {
  75. this.closed = true
  76. for (const waiters of this.waiters.values()) {
  77. for (const waiter of waiters) waiter.resolve()
  78. }
  79. this.waiters.clear()
  80. }
  81. }