rx.lua 20 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589590591592593594595596597598599600601602603604605606607608609610611612613614615616617618619620621622623624625626627628629630631632633634635636637638639640641642643644645646647648649650651652653654655656657658659660661662663664665666667668669670671672673674675676677678679680681682683684685686687688689690691692693694695696697698699700701702703704705706707708709710711712713714715716717718719720721722723724725726
  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 only produces the last result of the original.
  267. -- @returns {Observable}
  268. function Observable:last()
  269. return Observable.create(function(observer)
  270. local value
  271. local function onNext(...)
  272. value = {...}
  273. end
  274. local function onError(e)
  275. return observer:onError(e)
  276. end
  277. local function onComplete()
  278. observer:onNext(unpack(value or {}))
  279. return observer:onComplete()
  280. end
  281. return self:subscribe(onNext, onError, onComplete)
  282. end)
  283. end
  284. --- Returns a new Observable that produces the values of the original transformed by a function.
  285. -- @arg {function} callback - The function to transform values from the original Observable.
  286. -- @returns {Observable}
  287. function Observable:map(callback)
  288. return Observable.create(function(observer)
  289. callback = callback or identity
  290. local function onNext(...)
  291. return observer:onNext(callback(...))
  292. end
  293. local function onError(e)
  294. return observer:onError(e)
  295. end
  296. local function onComplete()
  297. return observer:onComplete()
  298. end
  299. return self:subscribe(onNext, onError, onComplete)
  300. end)
  301. end
  302. --- Returns a new Observable that produces the maximum value produced by the original.
  303. -- @returns {Observable}
  304. function Observable:max()
  305. return self:reduce(math.max)
  306. end
  307. --- Returns a new Observable that produces the minimum value produced by the original.
  308. -- @returns {Observable}
  309. function Observable:min()
  310. return self:reduce(math.min)
  311. end
  312. --- Returns a new Observable that produces the values produced by all the specified Observables in
  313. -- the order they are produced.
  314. -- @arg {Observable...} sources - One or more Observables to merge.
  315. -- @returns {Observable}
  316. function Observable:merge(...)
  317. local sources = {...}
  318. table.insert(sources, 1, self)
  319. return Observable.create(function(observer)
  320. local function onNext(...)
  321. return observer:onNext(...)
  322. end
  323. local function onError(message)
  324. return observer:onError(message)
  325. end
  326. local function onComplete(i)
  327. return function()
  328. sources[i] = nil
  329. if not next(sources) then
  330. observer:onComplete()
  331. end
  332. end
  333. end
  334. for i = 1, #sources do
  335. sources[i]:subscribe(onNext, onError, onComplete(i))
  336. end
  337. end)
  338. end
  339. --- Returns a new Observable that produces values computed by extracting the given key from the
  340. -- tables produced by the original.
  341. -- @arg {function} key - The key to extract from the table.
  342. -- @returns {Observable}
  343. function Observable:pluck(key)
  344. return Observable.create(function(observer)
  345. local function onNext(t)
  346. return observer:onNext(t[key])
  347. end
  348. local function onError(e)
  349. return observer:onError(e)
  350. end
  351. local function onComplete()
  352. return observer:onComplete()
  353. end
  354. return self:subscribe(onNext, onError, onComplete)
  355. end)
  356. end
  357. --- Returns a new Observable that produces a single value computed by accumulating the results of
  358. -- running a function on each value produced by the original Observable.
  359. -- @arg {function} accumulator - Accumulates the values of the original Observable. Will be passed
  360. -- the return value of the last call as the first argument and the
  361. -- current values as the rest of the arguments.
  362. -- @arg {*} seed - A value to pass to the accumulator the first time it is run.
  363. -- @returns {Observable}
  364. function Observable:reduce(accumulator, seed)
  365. return Observable.create(function(observer)
  366. local result
  367. local function onNext(...)
  368. result = result or seed or (...)
  369. result = accumulator(result, ...)
  370. end
  371. local function onError(e)
  372. return observer:onError(e)
  373. end
  374. local function onComplete()
  375. observer:onNext(result)
  376. return observer:onComplete()
  377. end
  378. return self:subscribe(onNext, onError, onComplete)
  379. end)
  380. end
  381. --- Returns a new Observable that skips over a specified number of values produced by the original
  382. -- and produces the rest.
  383. -- @arg {number=1} n - The number of values to ignore.
  384. -- @returns {Observable}
  385. function Observable:skip(n)
  386. n = n or 1
  387. return Observable.create(function(observer)
  388. local i = 1
  389. local function onNext(...)
  390. if i > n then
  391. observer:onNext(...)
  392. else
  393. i = i + 1
  394. end
  395. end
  396. local function onError(e)
  397. return observer:onError(e)
  398. end
  399. local function onComplete()
  400. return observer:onComplete()
  401. end
  402. return self:subscribe(onNext, onError, onComplete)
  403. end)
  404. end
  405. --- Returns a new Observable that skips over values produced by the original until the specified
  406. -- Observable produces a value.
  407. -- @arg {Observable} other - The Observable that triggers the production of values.
  408. -- @returns {Observable}
  409. function Observable:skipUntil(other)
  410. return Observable.create(function(observer)
  411. local function trigger()
  412. local function onNext(...)
  413. return observer:onNext(...)
  414. end
  415. local function onError(message)
  416. return observer:onNext(message)
  417. end
  418. local function onComplete()
  419. return observer:onComplete()
  420. end
  421. return self:subscribe(onNext, onError, onComplete)
  422. end
  423. other:subscribe(trigger, trigger, trigger)
  424. end)
  425. end
  426. --- Returns a new Observable that produces the sum of the values of the original Observable as a
  427. -- single result.
  428. -- @returns {Observable}
  429. function Observable:sum()
  430. return self:reduce(function(x, y) return x + y end, 0)
  431. end
  432. --- Returns a new Observable that only produces the first n results of the original.
  433. -- @arg {number=1} n - The number of elements to produce before completing.
  434. -- @returns {Observable}
  435. function Observable:take(n)
  436. n = n or 1
  437. return Observable.create(function(observer)
  438. if n <= 0 then
  439. observer:onComplete()
  440. return
  441. end
  442. local i = 1
  443. local function onNext(...)
  444. observer:onNext(...)
  445. i = i + 1
  446. if i > n then
  447. observer:onComplete()
  448. end
  449. end
  450. local function onError(e)
  451. return observer:onError(e)
  452. end
  453. local function onComplete()
  454. return observer:onComplete()
  455. end
  456. return self:subscribe(onNext, onError, onComplete)
  457. end)
  458. end
  459. --- Returns a new Observable that completes when the specified Observable fires.
  460. -- @arg {Observable} other - The Observable that triggers completion of the original.
  461. -- @returns {Observable}
  462. function Observable:takeUntil(other)
  463. return Observable.create(function(observer)
  464. local function onNext(...)
  465. return observer:onNext(...)
  466. end
  467. local function onError(e)
  468. return observer:onError(e)
  469. end
  470. local function onComplete()
  471. return observer:onComplete()
  472. end
  473. other:subscribe(onComplete, onComplete, onComplete)
  474. return self:subscribe(onNext, onError, onComplete)
  475. end)
  476. end
  477. --- @class Scheduler
  478. -- @description Schedulers manage groups of Observables.
  479. local Scheduler = {}
  480. --- @class CooperativeScheduler
  481. -- @description Manages Observables using coroutines and a virtual clock that must be updated
  482. -- manually.
  483. local Cooperative = {}
  484. Cooperative.__index = Cooperative
  485. --- Creates a new Cooperative Scheduler.
  486. -- @arg {number=0} currentTime - A time to start the scheduler at.
  487. -- @returns {Scheduler.Cooperative}
  488. function Cooperative.create(currentTime)
  489. local self = {
  490. tasks = {},
  491. currentTime = currentTime or 0
  492. }
  493. return setmetatable(self, Cooperative)
  494. end
  495. --- Schedules a function to be run after an optional delay.
  496. -- @arg {function} action - The function to execute. Will be converted into a coroutine. The
  497. -- coroutine may yield execution back to the scheduler with an optional
  498. -- number, which will put it to sleep for a time period.
  499. -- @arg {number=0} delay - Delay execution of the action by a time period.
  500. function Cooperative:schedule(action, delay)
  501. table.insert(self.tasks, {
  502. thread = coroutine.create(action),
  503. due = self.currentTime + (delay or 0)
  504. })
  505. end
  506. --- Triggers an update of the Cooperative Scheduler. The clock will be advanced and the scheduler
  507. -- will run any coroutines that are due to be run.
  508. -- @arg {number=0} delta - An amount of time to advance the clock by. It is common to pass in the
  509. -- time in seconds or milliseconds elapsed since this function was last
  510. -- called.
  511. function Cooperative:update(delta)
  512. self.currentTime = self.currentTime + (delta or 0)
  513. for i = #self.tasks, 1, -1 do
  514. local task = self.tasks[i]
  515. if self.currentTime >= task.due then
  516. local success, delay = coroutine.resume(task.thread)
  517. if success then
  518. task.due = math.max(task.due + (delay or 0), self.currentTime)
  519. else
  520. error(delay)
  521. end
  522. if coroutine.status(task.thread) == 'dead' then
  523. table.remove(self.tasks, i)
  524. end
  525. end
  526. end
  527. end
  528. --- Returns whether or not the Cooperative Scheduler's queue is empty.
  529. function Cooperative:isEmpty()
  530. return not next(self.tasks)
  531. end
  532. Scheduler.Cooperative = Cooperative
  533. --- @class Subject
  534. -- @description Subjects function both as an Observer and as an Observable. Subjects inherit all
  535. -- Observable functions, including subscribe. Values can also be pushed to the Subject, which will
  536. -- be broadcasted to any subscribed Observers.
  537. local Subject = setmetatable({}, Observable)
  538. Subject.__index = Subject
  539. --- Creates a new Subject.
  540. -- @arg {*...} value - The initial values.
  541. -- @returns {Subject}
  542. function Subject.create(...)
  543. local self = {
  544. value = {...},
  545. observers = {}
  546. }
  547. return setmetatable(self, Subject)
  548. end
  549. --- Creates a new Observer and attaches it to the Subject.
  550. -- @arg {function} onNext - Called when the Subject produces a value.
  551. -- @arg {function} onError - Called when the Subject terminates due to an error.
  552. -- @arg {function} onComplete - Called when the Subject completes normally.
  553. function Subject:subscribe(onNext, onError, onComplete)
  554. table.insert(self.observers, Observer.create(onNext, onError, onComplete))
  555. end
  556. --- Pushes zero or more values to the Subject. It will be broadcasted to all Observers.
  557. -- @arg {*...} values
  558. function Subject:onNext(...)
  559. self.value = {...}
  560. for i = 1, #self.observers do
  561. self.observers[i]:onNext(...)
  562. end
  563. end
  564. --- Signal to all Observers that an error has occurred.
  565. -- @arg {string=} message - A string describing what went wrong.
  566. function Subject:onError(message)
  567. for i = 1, #self.observers do
  568. self.observers[i]:onError(message)
  569. end
  570. end
  571. --- Signal to all Observers that the Subject will not produce any more values.
  572. function Subject:onComplete()
  573. for i = 1, #self.observers do
  574. self.observers[i]:onComplete()
  575. end
  576. end
  577. --- Returns the last value emitted by the Subject, or the initial value passed to the constructor
  578. -- if nothing has been emitted yet.
  579. -- @returns {*...}
  580. function Subject:getValue()
  581. return unpack(self.value or {})
  582. end
  583. Subject.__call = Subject.onNext
  584. rx = {
  585. Observer = Observer,
  586. Observable = Observable,
  587. Scheduler = Scheduler,
  588. scheduler = Scheduler.Cooperative.create(),
  589. Subject = Subject
  590. }
  591. return rx