rx.lua 20 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589590591592593594595596597598599600601602603604605606607608609610611612613614615616617618619620621622623624625626627628629630631632633634635636637638639640641642643644645646647648649650651652653654655656657658659660661662663664665666667668669670671672673674675676677678679680681682683684685686687688689690691692693694695696697698699700701702703704705706707708709710711712713714715716717718719720721722723724725726727728729730731732733734735736737738739740741742743744745746747748749750751
  1. local rx
  2. local function noop() end
  3. local function identity(x) return x end
  4. --- @class Observer
  5. -- @description Observers are simple objects that receive values from Observables.
  6. local Observer = {}
  7. Observer.__index = Observer
  8. --- Creates a new Observer.
  9. -- @arg {function=} onNext - Called when the Observable produces a value.
  10. -- @arg {function=} onError - Called when the Observable terminates due to an error.
  11. -- @arg {function=} onComplete - Called when the Observable completes normally.
  12. -- @returns {Observer}
  13. function Observer.create(onNext, onError, onComplete)
  14. local self = {
  15. _onNext = onNext or noop,
  16. _onError = onError or error,
  17. _onComplete = onComplete or noop,
  18. stopped = false
  19. }
  20. return setmetatable(self, Observer)
  21. end
  22. --- Pushes zero or more values to the Observer.
  23. -- @arg {*...} values
  24. function Observer:onNext(...)
  25. if not self.stopped then
  26. self._onNext(...)
  27. end
  28. end
  29. --- Notify the Observer that an error has occurred.
  30. -- @arg {string=} message - A string describing what went wrong.
  31. function Observer:onError(message)
  32. if not self.stopped then
  33. self.stopped = true
  34. self._onError(message)
  35. end
  36. end
  37. --- Notify the Observer that the sequence has completed and will produce no more values.
  38. function Observer:onComplete()
  39. if not self.stopped then
  40. self.stopped = true
  41. self._onComplete()
  42. end
  43. end
  44. --- @class Observable
  45. -- @description Observables push values to Observers.
  46. local Observable = {}
  47. Observable.__index = Observable
  48. --- Creates a new Observable.
  49. -- @arg {function} subscribe - The subscription function that produces values.
  50. -- @returns {Observable}
  51. function Observable.create(subscribe)
  52. local self = {
  53. _subscribe = subscribe
  54. }
  55. return setmetatable(self, Observable)
  56. end
  57. --- Creates an Observable that produces a single value.
  58. -- @arg {*} value
  59. -- @returns {Observable}
  60. function Observable.fromValue(value)
  61. return Observable.create(function(observer)
  62. observer:onNext(value)
  63. observer:onComplete()
  64. end)
  65. end
  66. --- Creates an Observable that produces a range of values in a manner similar to a Lua for loop.
  67. -- @arg {number} initial - The first value of the range, or the upper limit if no other arguments
  68. -- are specified.
  69. -- @arg {number=} limit - The second value of the range.
  70. -- @arg {number=1} step - An amount to increment the value by each iteration.
  71. -- @returns {Observable}
  72. function Observable.fromRange(initial, limit, step)
  73. if not limit and not step then
  74. initial, limit = 1, initial
  75. end
  76. step = step or 1
  77. return Observable.create(function(observer)
  78. for i = initial, limit, step do
  79. observer:onNext(i)
  80. end
  81. observer:onComplete()
  82. end)
  83. end
  84. --- Creates an Observable that produces values from a table.
  85. -- @arg {table} table - The table used to create the Observable.
  86. -- @arg {function=pairs} iterator - An iterator used to iterate the table, e.g. pairs or ipairs.
  87. -- @returns {Observable}
  88. function Observable.fromTable(t, iterator)
  89. iterator = iterator or pairs
  90. return Observable.create(function(observer)
  91. for key, value in iterator(t) do
  92. observer:onNext(value, key)
  93. end
  94. observer:onComplete()
  95. end)
  96. end
  97. --- Creates an Observable that produces values when the specified coroutine yields.
  98. -- @arg {thread} coroutine
  99. -- @returns {Observable}
  100. function Observable.fromCoroutine(thread)
  101. thread = type(thread) == 'function' and coroutine.create(thread) or thread
  102. return Observable.create(function(observer)
  103. return rx.scheduler:schedule(function()
  104. while not observer.stopped do
  105. local success, value = coroutine.resume(thread)
  106. if success then
  107. observer:onNext(value)
  108. else
  109. return observer:onError(value)
  110. end
  111. if coroutine.status(thread) == 'dead' then
  112. return observer:onComplete()
  113. end
  114. coroutine.yield()
  115. end
  116. end)
  117. end)
  118. end
  119. --- Shorthand for creating an Observer and passing it to this Observable's subscription function.
  120. -- @arg {function} onNext - Called when the Observable produces a value.
  121. -- @arg {function} onError - Called when the Observable terminates due to an error.
  122. -- @arg {function} onComplete - Called when the Observable completes normally.
  123. function Observable:subscribe(onNext, onError, onComplete)
  124. return self._subscribe(Observer.create(onNext, onError, onComplete))
  125. end
  126. --- Subscribes to this Observable and prints values it produces.
  127. -- @arg {string=} name - Prefixes the printed messages with a name.
  128. function Observable:dump(name)
  129. name = name and (name .. ' ') or ''
  130. local onNext = function(...) print(name .. 'onNext: ' .. table.concat({...}, ', ')) end
  131. local onError = function(e) print(name .. 'onError: ' .. e) end
  132. local onComplete = function() print(name .. 'onComplete') end
  133. return self:subscribe(onNext, onError, onComplete)
  134. end
  135. -- The functions below transform the values produced by an Observable and return a new Observable
  136. -- that produces these values.
  137. --- Returns a new Observable that runs a combinator function on the most recent values from a set
  138. -- of Observables whenever any of them produce a new value. The results of the combinator function
  139. -- are produced by the new Observable.
  140. -- @arg {Observable...} observables - One or more Observables to combine.
  141. -- @arg {function} combinator - A function that combines the latest result from each Observable and
  142. -- returns a single value.
  143. -- @returns {Observable}
  144. function Observable:combineLatest(...)
  145. local sources = {...}
  146. local combinator = table.remove(sources)
  147. table.insert(sources, 1, self)
  148. return Observable.create(function(observer)
  149. local latest = {}
  150. local pending = {unpack(sources)}
  151. local completed = {}
  152. local function onNext(i)
  153. return function(value)
  154. latest[i] = value
  155. pending[i] = nil
  156. if not next(pending) then
  157. observer:onNext(combinator(unpack(latest)))
  158. end
  159. end
  160. end
  161. local function onError(e)
  162. return observer:onError(e)
  163. end
  164. local function onComplete(i)
  165. return function()
  166. table.insert(completed, i)
  167. if #completed == #sources then
  168. observer:onComplete()
  169. end
  170. end
  171. end
  172. for i = 1, #sources do
  173. sources[i]:subscribe(onNext(i), onError, onComplete(i))
  174. end
  175. end)
  176. end
  177. --- Returns a new Observable that produces the values produced by all the specified Observables in
  178. -- the order they are specified.
  179. -- @arg {Observable...} sources - The Observables to concatenate.
  180. -- @returns {Observable}
  181. function Observable:concat(other, ...)
  182. if not other then return self end
  183. local others = {...}
  184. return Observable.create(function(observer)
  185. local function onNext(...)
  186. return observer:onNext(...)
  187. end
  188. local function onError(message)
  189. return observer:onError(message)
  190. end
  191. local function onComplete()
  192. return observer:onComplete()
  193. end
  194. local function chain()
  195. return other:concat(unpack(others)):subscribe(onNext, onError, onComplete)
  196. end
  197. return self:subscribe(onNext, onError, chain)
  198. end)
  199. end
  200. --- Returns a new Observable that produces the values from the original with duplicates removed.
  201. -- @returns {Observable}
  202. function Observable:distinct()
  203. return Observable.create(function(observer)
  204. local values = {}
  205. local function onNext(x)
  206. if not values[x] then
  207. observer:onNext(x)
  208. end
  209. values[x] = true
  210. end
  211. local function onError(e)
  212. return observer:onError(e)
  213. end
  214. local function onComplete()
  215. return observer:onComplete()
  216. end
  217. return self:subscribe(onNext, onError, onComplete)
  218. end)
  219. end
  220. --- Returns a new Observable that only produces values of the first that satisfy a predicate.
  221. -- @arg {function} predicate - The predicate used to filter values.
  222. -- @returns {Observable}
  223. function Observable:filter(predicate)
  224. predicate = predicate or identity
  225. return Observable.create(function(observer)
  226. local function onNext(...)
  227. if predicate(...) then
  228. return observer:onNext(...)
  229. end
  230. end
  231. local function onError(e)
  232. return observer:onError(e)
  233. end
  234. local function onComplete()
  235. return observer:onComplete(e)
  236. end
  237. return self:subscribe(onNext, onError, onComplete)
  238. end)
  239. end
  240. --- Returns a new Observable that produces the first value of the original that satisfies a
  241. -- predicate.
  242. -- @arg {function} predicate - The predicate used to find a value.
  243. function Observable:find(predicate)
  244. predicate = predicate or identity
  245. return Observable.create(function(observer)
  246. local function onNext(...)
  247. if predicate(...) then
  248. observer:onNext(...)
  249. return observer:onComplete()
  250. end
  251. end
  252. local function onError(message)
  253. return observer:onError(e)
  254. end
  255. local function onComplete()
  256. return observer:onComplete()
  257. end
  258. return self:subscribe(onNext, onError, onComplete)
  259. end)
  260. end
  261. --- Returns a new Observable that only produces the first result of the original.
  262. -- @returns {Observable}
  263. function Observable:first()
  264. return self:take(1)
  265. end
  266. --- Returns a new Observable that subscribes to the Observables produced by the original and
  267. -- produces their values.
  268. -- @returns {Observable}
  269. function Observable:flatten()
  270. return Observable.create(function(observer)
  271. local function onError(message)
  272. return observer:onError(message)
  273. end
  274. local function onNext(observable)
  275. local function innerOnNext(...)
  276. observer:onNext(...)
  277. end
  278. observable:subscribe(innerOnNext, onError, noop)
  279. end
  280. local function onComplete()
  281. return observer:onComplete()
  282. end
  283. return self:subscribe(onNext, onError, onComplete)
  284. end)
  285. end
  286. --- Returns a new Observable that only produces the last result of the original.
  287. -- @returns {Observable}
  288. function Observable:last()
  289. return Observable.create(function(observer)
  290. local value
  291. local function onNext(...)
  292. value = {...}
  293. end
  294. local function onError(e)
  295. return observer:onError(e)
  296. end
  297. local function onComplete()
  298. observer:onNext(unpack(value or {}))
  299. return observer:onComplete()
  300. end
  301. return self:subscribe(onNext, onError, onComplete)
  302. end)
  303. end
  304. --- Returns a new Observable that produces the values of the original transformed by a function.
  305. -- @arg {function} callback - The function to transform values from the original Observable.
  306. -- @returns {Observable}
  307. function Observable:map(callback)
  308. return Observable.create(function(observer)
  309. callback = callback or identity
  310. local function onNext(...)
  311. return observer:onNext(callback(...))
  312. end
  313. local function onError(e)
  314. return observer:onError(e)
  315. end
  316. local function onComplete()
  317. return observer:onComplete()
  318. end
  319. return self:subscribe(onNext, onError, onComplete)
  320. end)
  321. end
  322. --- Returns a new Observable that produces the maximum value produced by the original.
  323. -- @returns {Observable}
  324. function Observable:max()
  325. return self:reduce(math.max)
  326. end
  327. --- Returns a new Observable that produces the minimum value produced by the original.
  328. -- @returns {Observable}
  329. function Observable:min()
  330. return self:reduce(math.min)
  331. end
  332. --- Returns a new Observable that produces the values produced by all the specified Observables in
  333. -- the order they are produced.
  334. -- @arg {Observable...} sources - One or more Observables to merge.
  335. -- @returns {Observable}
  336. function Observable:merge(...)
  337. local sources = {...}
  338. table.insert(sources, 1, self)
  339. return Observable.create(function(observer)
  340. local function onNext(...)
  341. return observer:onNext(...)
  342. end
  343. local function onError(message)
  344. return observer:onError(message)
  345. end
  346. local function onComplete(i)
  347. return function()
  348. sources[i] = nil
  349. if not next(sources) then
  350. observer:onComplete()
  351. end
  352. end
  353. end
  354. for i = 1, #sources do
  355. sources[i]:subscribe(onNext, onError, onComplete(i))
  356. end
  357. end)
  358. end
  359. --- Returns a new Observable that produces values computed by extracting the given key from the
  360. -- tables produced by the original.
  361. -- @arg {function} key - The key to extract from the table.
  362. -- @returns {Observable}
  363. function Observable:pluck(key)
  364. return Observable.create(function(observer)
  365. local function onNext(t)
  366. return observer:onNext(t[key])
  367. end
  368. local function onError(e)
  369. return observer:onError(e)
  370. end
  371. local function onComplete()
  372. return observer:onComplete()
  373. end
  374. return self:subscribe(onNext, onError, onComplete)
  375. end)
  376. end
  377. --- Returns a new Observable that produces a single value computed by accumulating the results of
  378. -- running a function on each value produced by the original Observable.
  379. -- @arg {function} accumulator - Accumulates the values of the original Observable. Will be passed
  380. -- the return value of the last call as the first argument and the
  381. -- current values as the rest of the arguments.
  382. -- @arg {*} seed - A value to pass to the accumulator the first time it is run.
  383. -- @returns {Observable}
  384. function Observable:reduce(accumulator, seed)
  385. return Observable.create(function(observer)
  386. local result
  387. local function onNext(...)
  388. result = result or seed or (...)
  389. result = accumulator(result, ...)
  390. end
  391. local function onError(e)
  392. return observer:onError(e)
  393. end
  394. local function onComplete()
  395. observer:onNext(result)
  396. return observer:onComplete()
  397. end
  398. return self:subscribe(onNext, onError, onComplete)
  399. end)
  400. end
  401. --- Returns a new Observable that skips over a specified number of values produced by the original
  402. -- and produces the rest.
  403. -- @arg {number=1} n - The number of values to ignore.
  404. -- @returns {Observable}
  405. function Observable:skip(n)
  406. n = n or 1
  407. return Observable.create(function(observer)
  408. local i = 1
  409. local function onNext(...)
  410. if i > n then
  411. observer:onNext(...)
  412. else
  413. i = i + 1
  414. end
  415. end
  416. local function onError(e)
  417. return observer:onError(e)
  418. end
  419. local function onComplete()
  420. return observer:onComplete()
  421. end
  422. return self:subscribe(onNext, onError, onComplete)
  423. end)
  424. end
  425. --- Returns a new Observable that skips over values produced by the original until the specified
  426. -- Observable produces a value.
  427. -- @arg {Observable} other - The Observable that triggers the production of values.
  428. -- @returns {Observable}
  429. function Observable:skipUntil(other)
  430. return Observable.create(function(observer)
  431. local function trigger()
  432. local function onNext(...)
  433. return observer:onNext(...)
  434. end
  435. local function onError(message)
  436. return observer:onNext(message)
  437. end
  438. local function onComplete()
  439. return observer:onComplete()
  440. end
  441. return self:subscribe(onNext, onError, onComplete)
  442. end
  443. other:subscribe(trigger, trigger, trigger)
  444. end)
  445. end
  446. --- Returns a new Observable that produces the sum of the values of the original Observable as a
  447. -- single result.
  448. -- @returns {Observable}
  449. function Observable:sum()
  450. return self:reduce(function(x, y) return x + y end, 0)
  451. end
  452. --- Returns a new Observable that only produces the first n results of the original.
  453. -- @arg {number=1} n - The number of elements to produce before completing.
  454. -- @returns {Observable}
  455. function Observable:take(n)
  456. n = n or 1
  457. return Observable.create(function(observer)
  458. if n <= 0 then
  459. observer:onComplete()
  460. return
  461. end
  462. local i = 1
  463. local function onNext(...)
  464. observer:onNext(...)
  465. i = i + 1
  466. if i > n then
  467. observer:onComplete()
  468. end
  469. end
  470. local function onError(e)
  471. return observer:onError(e)
  472. end
  473. local function onComplete()
  474. return observer:onComplete()
  475. end
  476. return self:subscribe(onNext, onError, onComplete)
  477. end)
  478. end
  479. --- Returns a new Observable that completes when the specified Observable fires.
  480. -- @arg {Observable} other - The Observable that triggers completion of the original.
  481. -- @returns {Observable}
  482. function Observable:takeUntil(other)
  483. return Observable.create(function(observer)
  484. local function onNext(...)
  485. return observer:onNext(...)
  486. end
  487. local function onError(e)
  488. return observer:onError(e)
  489. end
  490. local function onComplete()
  491. return observer:onComplete()
  492. end
  493. other:subscribe(onComplete, onComplete, onComplete)
  494. return self:subscribe(onNext, onError, onComplete)
  495. end)
  496. end
  497. --- @class Scheduler
  498. -- @description Schedulers manage groups of Observables.
  499. local Scheduler = {}
  500. --- @class CooperativeScheduler
  501. -- @description Manages Observables using coroutines and a virtual clock that must be updated
  502. -- manually.
  503. local Cooperative = {}
  504. Cooperative.__index = Cooperative
  505. --- Creates a new Cooperative Scheduler.
  506. -- @arg {number=0} currentTime - A time to start the scheduler at.
  507. -- @returns {Scheduler.Cooperative}
  508. function Cooperative.create(currentTime)
  509. local self = {
  510. tasks = {},
  511. currentTime = currentTime or 0
  512. }
  513. return setmetatable(self, Cooperative)
  514. end
  515. --- Schedules a function to be run after an optional delay.
  516. -- @arg {function} action - The function to execute. Will be converted into a coroutine. The
  517. -- coroutine may yield execution back to the scheduler with an optional
  518. -- number, which will put it to sleep for a time period.
  519. -- @arg {number=0} delay - Delay execution of the action by a time period.
  520. function Cooperative:schedule(action, delay)
  521. table.insert(self.tasks, {
  522. thread = coroutine.create(action),
  523. due = self.currentTime + (delay or 0)
  524. })
  525. end
  526. --- Triggers an update of the Cooperative Scheduler. The clock will be advanced and the scheduler
  527. -- will run any coroutines that are due to be run.
  528. -- @arg {number=0} delta - An amount of time to advance the clock by. It is common to pass in the
  529. -- time in seconds or milliseconds elapsed since this function was last
  530. -- called.
  531. function Cooperative:update(delta)
  532. self.currentTime = self.currentTime + (delta or 0)
  533. for i = #self.tasks, 1, -1 do
  534. local task = self.tasks[i]
  535. if self.currentTime >= task.due then
  536. local success, delay = coroutine.resume(task.thread)
  537. if success then
  538. task.due = math.max(task.due + (delay or 0), self.currentTime)
  539. else
  540. error(delay)
  541. end
  542. if coroutine.status(task.thread) == 'dead' then
  543. table.remove(self.tasks, i)
  544. end
  545. end
  546. end
  547. end
  548. --- Returns whether or not the Cooperative Scheduler's queue is empty.
  549. function Cooperative:isEmpty()
  550. return not next(self.tasks)
  551. end
  552. Scheduler.Cooperative = Cooperative
  553. --- @class Subject
  554. -- @description Subjects function both as an Observer and as an Observable. Subjects inherit all
  555. -- Observable functions, including subscribe. Values can also be pushed to the Subject, which will
  556. -- be broadcasted to any subscribed Observers.
  557. local Subject = setmetatable({}, Observable)
  558. Subject.__index = Subject
  559. --- Creates a new Subject.
  560. -- @arg {*...} value - The initial values.
  561. -- @returns {Subject}
  562. function Subject.create(...)
  563. local self = {
  564. value = {...},
  565. observers = {}
  566. }
  567. return setmetatable(self, Subject)
  568. end
  569. --- Creates a new Observer and attaches it to the Subject.
  570. -- @arg {function} onNext - Called when the Subject produces a value.
  571. -- @arg {function} onError - Called when the Subject terminates due to an error.
  572. -- @arg {function} onComplete - Called when the Subject completes normally.
  573. function Subject:subscribe(onNext, onError, onComplete)
  574. table.insert(self.observers, Observer.create(onNext, onError, onComplete))
  575. end
  576. --- Pushes zero or more values to the Subject. It will be broadcasted to all Observers.
  577. -- @arg {*...} values
  578. function Subject:onNext(...)
  579. self.value = {...}
  580. for i = 1, #self.observers do
  581. self.observers[i]:onNext(...)
  582. end
  583. end
  584. --- Signal to all Observers that an error has occurred.
  585. -- @arg {string=} message - A string describing what went wrong.
  586. function Subject:onError(message)
  587. for i = 1, #self.observers do
  588. self.observers[i]:onError(message)
  589. end
  590. end
  591. --- Signal to all Observers that the Subject will not produce any more values.
  592. function Subject:onComplete()
  593. for i = 1, #self.observers do
  594. self.observers[i]:onComplete()
  595. end
  596. end
  597. --- Returns the last value emitted by the Subject, or the initial value passed to the constructor
  598. -- if nothing has been emitted yet.
  599. -- @returns {*...}
  600. function Subject:getValue()
  601. return unpack(self.value or {})
  602. end
  603. Subject.__call = Subject.onNext
  604. rx = {
  605. Observer = Observer,
  606. Observable = Observable,
  607. Scheduler = Scheduler,
  608. scheduler = Scheduler.Cooperative.create(),
  609. Subject = Subject
  610. }
  611. return rx