Flatten.swift 36 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589590591592593594595596597598599600601602603604605606607608609610611612613614615616617618619620621622623624625626627628629630631632633634635636637638639640641642643644645646647648649650651652653654655656657658659660661662663664665666667668669670671672673674675676677678679680681682683684685686687688689690691692693694695696697698699700701702703704705706707708709710711712713714715716717718719720721722723724725726727728729730731732733734735736737738739740741742743744745746747748749750751752753754755756757758759760761762763764765766767768769770771772773774775776777778779780781782783784785786787788789790791792793794795796797798799800801802803804805806807808809810811812813814815816817818819820821822823824825826827828829830831832833834835836837838839840841842843844845846847848849850851852853854855856857858859860861862863864865866867868869870871872873874875876877878879880881882883884885886887888889890891892893894895896897898899900901902903904905906907908909910911912913914915916917918919920921922923924925926927928929930931932933934935936937938939940941942943944945946947948949950951952953954955956957958959960961962963964965966967968969970971972973974
  1. //
  2. // Flatten.swift
  3. // ReactiveCocoa
  4. //
  5. // Created by Neil Pankey on 11/30/15.
  6. // Copyright © 2015 GitHub. All rights reserved.
  7. //
  8. import enum Result.NoError
  9. /// Describes how multiple producers should be joined together.
  10. public enum FlattenStrategy: Equatable {
  11. /// The producers should be merged, so that any value received on any of the
  12. /// input producers will be forwarded immediately to the output producer.
  13. ///
  14. /// The resulting producer will complete only when all inputs have
  15. /// completed.
  16. case Merge
  17. /// The producers should be concatenated, so that their values are sent in
  18. /// the order of the producers themselves.
  19. ///
  20. /// The resulting producer will complete only when all inputs have
  21. /// completed.
  22. case Concat
  23. /// Only the events from the latest input producer should be considered for
  24. /// the output. Any producers received before that point will be disposed
  25. /// of.
  26. ///
  27. /// The resulting producer will complete only when the producer-of-producers
  28. /// and the latest producer has completed.
  29. case Latest
  30. }
  31. extension SignalType where Value: SignalProducerType, Error == Value.Error {
  32. /// Flattens the inner producers sent upon `signal` (into a single signal of
  33. /// values), according to the semantics of the given strategy.
  34. ///
  35. /// - note: If `signal` or an active inner producer fails, the returned
  36. /// signal will forward that failure immediately.
  37. ///
  38. /// - note: `Interrupted` events on inner producers will be treated like
  39. /// `Completed events on inner producers.
  40. @warn_unused_result(message="Did you forget to call `observe` on the signal?")
  41. public func flatten(strategy: FlattenStrategy) -> Signal<Value.Value, Error> {
  42. switch strategy {
  43. case .Merge:
  44. return self.merge()
  45. case .Concat:
  46. return self.concat()
  47. case .Latest:
  48. return self.switchToLatest()
  49. }
  50. }
  51. }
  52. extension SignalType where Value: SignalProducerType, Error == NoError {
  53. /// Flattens the inner producers sent upon `signal` (into a single signal of
  54. /// values), according to the semantics of the given strategy.
  55. ///
  56. /// - note: If an active inner producer fails, the returned signal will
  57. /// forward that failure immediately.
  58. ///
  59. /// - warning: `Interrupted` events on inner producers will be treated like
  60. /// `Completed` events on inner producers.
  61. ///
  62. /// - parameters:
  63. /// - strategy: Strategy used when flattening signals.
  64. @warn_unused_result(message="Did you forget to call `observe` on the signal?")
  65. public func flatten(strategy: FlattenStrategy) -> Signal<Value.Value, Value.Error> {
  66. return self
  67. .promoteErrors(Value.Error.self)
  68. .flatten(strategy)
  69. }
  70. }
  71. extension SignalType where Value: SignalProducerType, Error == NoError, Value.Error == NoError {
  72. /// Flattens the inner producers sent upon `signal` (into a single signal of
  73. /// values), according to the semantics of the given strategy.
  74. ///
  75. /// - warning: `Interrupted` events on inner producers will be treated like
  76. /// `Completed` events on inner producers.
  77. ///
  78. /// - parameters:
  79. /// - strategy: Strategy used when flattening signals.
  80. @warn_unused_result(message="Did you forget to call `observe` on the signal?")
  81. public func flatten(strategy: FlattenStrategy) -> Signal<Value.Value, Value.Error> {
  82. switch strategy {
  83. case .Merge:
  84. return self.merge()
  85. case .Concat:
  86. return self.concat()
  87. case .Latest:
  88. return self.switchToLatest()
  89. }
  90. }
  91. }
  92. extension SignalType where Value: SignalProducerType, Value.Error == NoError {
  93. /// Flattens the inner producers sent upon `signal` (into a single signal of
  94. /// values), according to the semantics of the given strategy.
  95. ///
  96. /// - note: If `signal` fails, the returned signal will forward that failure
  97. /// immediately.
  98. ///
  99. /// - warning: `Interrupted` events on inner producers will be treated like
  100. /// `Completed` events on inner producers.
  101. @warn_unused_result(message="Did you forget to call `observe` on the signal?")
  102. public func flatten(strategy: FlattenStrategy) -> Signal<Value.Value, Error> {
  103. return self.flatMap(strategy) { $0.promoteErrors(Error.self) }
  104. }
  105. }
  106. extension SignalProducerType where Value: SignalProducerType, Error == Value.Error {
  107. /// Flattens the inner producers sent upon `producer` (into a single
  108. /// producer of values), according to the semantics of the given strategy.
  109. ///
  110. /// - note: If `producer` or an active inner producer fails, the returned
  111. /// producer will forward that failure immediately.
  112. ///
  113. /// - warning: `Interrupted` events on inner producers will be treated like
  114. /// `Completed` events on inner producers.
  115. @warn_unused_result(message="Did you forget to call `start` on the producer?")
  116. public func flatten(strategy: FlattenStrategy) -> SignalProducer<Value.Value, Error> {
  117. switch strategy {
  118. case .Merge:
  119. return self.merge()
  120. case .Concat:
  121. return self.concat()
  122. case .Latest:
  123. return self.switchToLatest()
  124. }
  125. }
  126. }
  127. extension SignalProducerType where Value: SignalProducerType, Error == NoError {
  128. /// Flattens the inner producers sent upon `producer` (into a single
  129. /// producer of values), according to the semantics of the given strategy.
  130. ///
  131. /// - note: If an active inner producer fails, the returned producer will
  132. /// forward that failure immediately.
  133. ///
  134. /// - warning: `Interrupted` events on inner producers will be treated like
  135. /// `Completed` events on inner producers.
  136. @warn_unused_result(message="Did you forget to call `start` on the producer?")
  137. public func flatten(strategy: FlattenStrategy) -> SignalProducer<Value.Value, Value.Error> {
  138. return self
  139. .promoteErrors(Value.Error.self)
  140. .flatten(strategy)
  141. }
  142. }
  143. extension SignalProducerType where Value: SignalProducerType, Error == NoError, Value.Error == NoError {
  144. /// Flattens the inner producers sent upon `producer` (into a single
  145. /// producer of values), according to the semantics of the given strategy.
  146. ///
  147. /// - warning: `Interrupted` events on inner producers will be treated like
  148. /// `Completed` events on inner producers.
  149. @warn_unused_result(message="Did you forget to call `observe` on the signal?")
  150. public func flatten(strategy: FlattenStrategy) -> SignalProducer<Value.Value, Value.Error> {
  151. switch strategy {
  152. case .Merge:
  153. return self.merge()
  154. case .Concat:
  155. return self.concat()
  156. case .Latest:
  157. return self.switchToLatest()
  158. }
  159. }
  160. }
  161. extension SignalProducerType where Value: SignalProducerType, Value.Error == NoError {
  162. /// Flattens the inner producers sent upon `signal` (into a single signal of
  163. /// values), according to the semantics of the given strategy.
  164. ///
  165. /// - note: If `signal` fails, the returned signal will forward that failure
  166. /// immediately.
  167. ///
  168. /// - warning: `Interrupted` events on inner producers will be treated like
  169. /// `Completed` events on inner producers.
  170. @warn_unused_result(message="Did you forget to call `observe` on the signal?")
  171. public func flatten(strategy: FlattenStrategy) -> SignalProducer<Value.Value, Error> {
  172. return self.flatMap(strategy) { $0.promoteErrors(Error.self) }
  173. }
  174. }
  175. extension SignalType where Value: SignalType, Error == Value.Error {
  176. /// Flattens the inner signals sent upon `signal` (into a single signal of
  177. /// values), according to the semantics of the given strategy.
  178. ///
  179. /// - note: If `signal` or an active inner signal emits an error, the
  180. /// returned signal will forward that error immediately.
  181. ///
  182. /// - warning: `Interrupted` events on inner signals will be treated like
  183. /// `Completed` events on inner signals.
  184. @warn_unused_result(message="Did you forget to call `observe` on the signal?")
  185. public func flatten(strategy: FlattenStrategy) -> Signal<Value.Value, Error> {
  186. return self
  187. .map(SignalProducer.init)
  188. .flatten(strategy)
  189. }
  190. }
  191. extension SignalType where Value: SignalType, Error == NoError {
  192. /// Flattens the inner signals sent upon `signal` (into a single signal of
  193. /// values), according to the semantics of the given strategy.
  194. ///
  195. /// - note: If an active inner signal emits an error, the returned signal
  196. /// will forward that error immediately.
  197. ///
  198. /// - warning: `Interrupted` events on inner signals will be treated like
  199. /// `Completed` events on inner signals.
  200. @warn_unused_result(message="Did you forget to call `observe` on the signal?")
  201. public func flatten(strategy: FlattenStrategy) -> Signal<Value.Value, Value.Error> {
  202. return self
  203. .promoteErrors(Value.Error.self)
  204. .flatten(strategy)
  205. }
  206. }
  207. extension SignalType where Value: SignalType, Error == NoError, Value.Error == NoError {
  208. /// Flattens the inner signals sent upon `signal` (into a single signal of
  209. /// values), according to the semantics of the given strategy.
  210. ///
  211. /// - warning: `Interrupted` events on inner signals will be treated like
  212. /// `Completed` events on inner signals.
  213. @warn_unused_result(message="Did you forget to call `observe` on the signal?")
  214. public func flatten(strategy: FlattenStrategy) -> Signal<Value.Value, Value.Error> {
  215. return self
  216. .map(SignalProducer.init)
  217. .flatten(strategy)
  218. }
  219. }
  220. extension SignalType where Value: SignalType, Value.Error == NoError {
  221. /// Flattens the inner signals sent upon `signal` (into a single signal of
  222. /// values), according to the semantics of the given strategy.
  223. ///
  224. /// - note: If `signal` emits an error, the returned signal will forward
  225. /// that error immediately.
  226. ///
  227. /// - warning: `Interrupted` events on inner signals will be treated like
  228. /// `Completed` events on inner signals.
  229. @warn_unused_result(message="Did you forget to call `observe` on the signal?")
  230. public func flatten(strategy: FlattenStrategy) -> Signal<Value.Value, Error> {
  231. return self.flatMap(strategy) { $0.promoteErrors(Error.self) }
  232. }
  233. }
  234. extension SignalType where Value: SequenceType, Error == NoError {
  235. /// Flattens the `sequence` value sent by `signal` according to
  236. /// the semantics of the given strategy.
  237. @warn_unused_result(message="Did you forget to call `observe` on the signal?")
  238. public func flatten(strategy: FlattenStrategy) -> Signal<Value.Generator.Element, Error> {
  239. return self.flatMap(strategy) { .init(values: $0) }
  240. }
  241. }
  242. extension SignalProducerType where Value: SignalType, Error == Value.Error {
  243. /// Flattens the inner signals sent upon `producer` (into a single producer
  244. /// of values), according to the semantics of the given strategy.
  245. ///
  246. /// - note: If `producer` or an active inner signal emits an error, the
  247. /// returned producer will forward that error immediately.
  248. ///
  249. /// - warning: `Interrupted` events on inner signals will be treated like
  250. /// `Completed` events on inner signals.
  251. @warn_unused_result(message="Did you forget to call `start` on the producer?")
  252. public func flatten(strategy: FlattenStrategy) -> SignalProducer<Value.Value, Error> {
  253. return self
  254. .map(SignalProducer.init)
  255. .flatten(strategy)
  256. }
  257. }
  258. extension SignalProducerType where Value: SignalType, Error == NoError {
  259. /// Flattens the inner signals sent upon `producer` (into a single producer
  260. /// of values), according to the semantics of the given strategy.
  261. ///
  262. /// - note: If an active inner signal emits an error, the returned producer
  263. /// will forward that error immediately.
  264. ///
  265. /// - warning: `Interrupted` events on inner signals will be treated like
  266. /// `Completed` events on inner signals.
  267. @warn_unused_result(message="Did you forget to call `observe` on the signal?")
  268. public func flatten(strategy: FlattenStrategy) -> SignalProducer<Value.Value, Value.Error> {
  269. return self
  270. .promoteErrors(Value.Error.self)
  271. .flatten(strategy)
  272. }
  273. }
  274. extension SignalProducerType where Value: SignalType, Error == NoError, Value.Error == NoError {
  275. /// Flattens the inner signals sent upon `producer` (into a single producer
  276. /// of values), according to the semantics of the given strategy.
  277. ///
  278. /// - warning: `Interrupted` events on inner signals will be treated like
  279. /// `Completed` events on inner signals.
  280. @warn_unused_result(message="Did you forget to call `observe` on the signal?")
  281. public func flatten(strategy: FlattenStrategy) -> SignalProducer<Value.Value, Value.Error> {
  282. return self
  283. .map(SignalProducer.init)
  284. .flatten(strategy)
  285. }
  286. }
  287. extension SignalProducerType where Value: SignalType, Value.Error == NoError {
  288. /// Flattens the inner signals sent upon `producer` (into a single producer
  289. /// of values), according to the semantics of the given strategy.
  290. ///
  291. /// - note: If `producer` emits an error, the returned producer will forward
  292. /// that error immediately.
  293. ///
  294. /// - warning: `Interrupted` events on inner signals will be treated like
  295. /// `Completed` events on inner signals.
  296. @warn_unused_result(message="Did you forget to call `observe` on the signal?")
  297. public func flatten(strategy: FlattenStrategy) -> SignalProducer<Value.Value, Error> {
  298. return self.flatMap(strategy) { $0.promoteErrors(Error.self) }
  299. }
  300. }
  301. extension SignalProducerType where Value: SequenceType, Error == NoError {
  302. /// Flattens the `sequence` value sent by `producer` according to
  303. /// the semantics of the given strategy.
  304. @warn_unused_result(message="Did you forget to call `observe` on the signal?")
  305. public func flatten(strategy: FlattenStrategy) -> SignalProducer<Value.Generator.Element, Error> {
  306. return self.flatMap(strategy) { .init(values: $0) }
  307. }
  308. }
  309. extension SignalType where Value: SignalProducerType, Error == Value.Error {
  310. /// Returns a signal which sends all the values from producer signal emitted
  311. /// from `signal`, waiting until each inner producer completes before
  312. /// beginning to send the values from the next inner producer.
  313. ///
  314. /// - note: If any of the inner producers fail, the returned signal will
  315. /// forward that failure immediately
  316. ///
  317. /// - note: The returned signal completes only when `signal` and all
  318. /// producers emitted from `signal` complete.
  319. private func concat() -> Signal<Value.Value, Error> {
  320. return Signal<Value.Value, Error> { relayObserver in
  321. let disposable = CompositeDisposable()
  322. let relayDisposable = CompositeDisposable()
  323. disposable += relayDisposable
  324. disposable += self.observeConcat(relayObserver, relayDisposable)
  325. return disposable
  326. }
  327. }
  328. private func observeConcat(observer: Observer<Value.Value, Error>, _ disposable: CompositeDisposable? = nil) -> Disposable? {
  329. let state = ConcatState(observer: observer, disposable: disposable)
  330. return self.observe { event in
  331. switch event {
  332. case let .Next(value):
  333. state.enqueueSignalProducer(value.producer)
  334. case let .Failed(error):
  335. observer.sendFailed(error)
  336. case .Completed:
  337. // Add one last producer to the queue, whose sole job is to
  338. // "turn out the lights" by completing `observer`.
  339. state.enqueueSignalProducer(SignalProducer.empty.on(completed: {
  340. observer.sendCompleted()
  341. }))
  342. case .Interrupted:
  343. observer.sendInterrupted()
  344. }
  345. }
  346. }
  347. }
  348. extension SignalProducerType where Value: SignalProducerType, Error == Value.Error {
  349. /// Returns a producer which sends all the values from each producer emitted
  350. /// from `producer`, waiting until each inner producer completes before
  351. /// beginning to send the values from the next inner producer.
  352. ///
  353. /// - note: If any of the inner producers emit an error, the returned
  354. /// producer will emit that error.
  355. ///
  356. /// - note: The returned producer completes only when `producer` and all
  357. /// producers emitted from `producer` complete.
  358. private func concat() -> SignalProducer<Value.Value, Error> {
  359. return SignalProducer<Value.Value, Error> { observer, disposable in
  360. self.startWithSignal { signal, signalDisposable in
  361. disposable += signalDisposable
  362. signal.observeConcat(observer, disposable)
  363. }
  364. }
  365. }
  366. }
  367. extension SignalProducerType {
  368. /// `concat`s `next` onto `self`.
  369. @warn_unused_result(message="Did you forget to call `start` on the producer?")
  370. public func concat(next: SignalProducer<Value, Error>) -> SignalProducer<Value, Error> {
  371. return SignalProducer<SignalProducer<Value, Error>, Error>(values: [ self.producer, next ]).flatten(.Concat)
  372. }
  373. /// `concat`s `value` onto `self`.
  374. @warn_unused_result(message="Did you forget to call `start` on the producer?")
  375. public func concat(value value: Value) -> SignalProducer<Value, Error> {
  376. return self.concat(SignalProducer(value: value))
  377. }
  378. /// `concat`s `self` onto initial `previous`.
  379. @warn_unused_result(message="Did you forget to call `start` on the producer?")
  380. public func prefix<P: SignalProducerType where P.Value == Value, P.Error == Error>(previous: P) -> SignalProducer<Value, Error> {
  381. return previous.concat(self.producer)
  382. }
  383. /// `concat`s `self` onto initial `value`.
  384. @warn_unused_result(message="Did you forget to call `start` on the producer?")
  385. public func prefix(value value: Value) -> SignalProducer<Value, Error> {
  386. return self.prefix(SignalProducer(value: value))
  387. }
  388. }
  389. private final class ConcatState<Value, Error: ErrorType> {
  390. /// The observer of a started `concat` producer.
  391. let observer: Observer<Value, Error>
  392. /// The top level disposable of a started `concat` producer.
  393. let disposable: CompositeDisposable?
  394. /// The active producer, if any, and the producers waiting to be started.
  395. let queuedSignalProducers: Atomic<[SignalProducer<Value, Error>]> = Atomic([])
  396. init(observer: Signal<Value, Error>.Observer, disposable: CompositeDisposable?) {
  397. self.observer = observer
  398. self.disposable = disposable
  399. }
  400. func enqueueSignalProducer(producer: SignalProducer<Value, Error>) {
  401. if let d = disposable where d.disposed {
  402. return
  403. }
  404. var shouldStart = true
  405. queuedSignalProducers.modify {
  406. // An empty queue means the concat is idle, ready & waiting to start
  407. // the next producer.
  408. var queue = $0
  409. shouldStart = queue.isEmpty
  410. queue.append(producer)
  411. return queue
  412. }
  413. if shouldStart {
  414. startNextSignalProducer(producer)
  415. }
  416. }
  417. func dequeueSignalProducer() -> SignalProducer<Value, Error>? {
  418. if let d = disposable where d.disposed {
  419. return nil
  420. }
  421. var nextSignalProducer: SignalProducer<Value, Error>?
  422. queuedSignalProducers.modify {
  423. // Active producers remain in the queue until completed. Since
  424. // dequeueing happens at completion of the active producer, the
  425. // first producer in the queue can be removed.
  426. var queue = $0
  427. if !queue.isEmpty { queue.removeAtIndex(0) }
  428. nextSignalProducer = queue.first
  429. return queue
  430. }
  431. return nextSignalProducer
  432. }
  433. /// Subscribes to the given signal producer.
  434. func startNextSignalProducer(signalProducer: SignalProducer<Value, Error>) {
  435. signalProducer.startWithSignal { signal, disposable in
  436. let handle = self.disposable?.addDisposable(disposable) ?? nil
  437. signal.observe { event in
  438. switch event {
  439. case .Completed, .Interrupted:
  440. handle?.remove()
  441. if let nextSignalProducer = self.dequeueSignalProducer() {
  442. self.startNextSignalProducer(nextSignalProducer)
  443. }
  444. case .Next, .Failed:
  445. self.observer.action(event)
  446. }
  447. }
  448. }
  449. }
  450. }
  451. extension SignalType where Value: SignalProducerType, Error == Value.Error {
  452. /// Merges a `signal` of SignalProducers down into a single signal, biased
  453. /// toward the producer added earlier. Returns a Signal that will forward
  454. /// events from the inner producers as they arrive.
  455. private func merge() -> Signal<Value.Value, Error> {
  456. return Signal<Value.Value, Error> { relayObserver in
  457. let disposable = CompositeDisposable()
  458. let relayDisposable = CompositeDisposable()
  459. disposable += relayDisposable
  460. disposable += self.observeMerge(relayObserver, relayDisposable)
  461. return disposable
  462. }
  463. }
  464. private func observeMerge(observer: Observer<Value.Value, Error>, _ disposable: CompositeDisposable) -> Disposable? {
  465. let inFlight = Atomic(1)
  466. let decrementInFlight = {
  467. let orig = inFlight.modify { $0 - 1 }
  468. if orig == 1 {
  469. observer.sendCompleted()
  470. }
  471. }
  472. return self.observe { event in
  473. switch event {
  474. case let .Next(producer):
  475. producer.startWithSignal { innerSignal, innerDisposable in
  476. inFlight.modify { $0 + 1 }
  477. let handle = disposable.addDisposable(innerDisposable)
  478. innerSignal.observe { event in
  479. switch event {
  480. case .Completed, .Interrupted:
  481. handle.remove()
  482. decrementInFlight()
  483. case .Next, .Failed:
  484. observer.action(event)
  485. }
  486. }
  487. }
  488. case let .Failed(error):
  489. observer.sendFailed(error)
  490. case .Completed:
  491. decrementInFlight()
  492. case .Interrupted:
  493. observer.sendInterrupted()
  494. }
  495. }
  496. }
  497. }
  498. extension SignalProducerType where Value: SignalProducerType, Error == Value.Error {
  499. /// Merges a `signal` of SignalProducers down into a single signal, biased
  500. /// toward the producer added earlier. Returns a Signal that will forward
  501. /// events from the inner producers as they arrive.
  502. private func merge() -> SignalProducer<Value.Value, Error> {
  503. return SignalProducer<Value.Value, Error> { relayObserver, disposable in
  504. self.startWithSignal { signal, signalDisposable in
  505. disposable.addDisposable(signalDisposable)
  506. signal.observeMerge(relayObserver, disposable)
  507. }
  508. }
  509. }
  510. }
  511. extension SignalType {
  512. /// Merges the given signals into a single `Signal` that will emit all
  513. /// values from each of them, and complete when all of them have completed.
  514. @warn_unused_result(message="Did you forget to call `observe` on the signal?")
  515. public static func merge<Seq: SequenceType, S: SignalType where S.Value == Value, S.Error == Error, Seq.Generator.Element == S>(signals: Seq) -> Signal<Value, Error> {
  516. let producer = SignalProducer<S, Error>(values: signals)
  517. var result: Signal<Value, Error>!
  518. producer.startWithSignal { signal, _ in
  519. result = signal.flatten(.Merge)
  520. }
  521. return result
  522. }
  523. /// Merges the given signals into a single `Signal` that will emit all
  524. /// values from each of them, and complete when all of them have completed.
  525. @warn_unused_result(message="Did you forget to call `observe` on the signal?")
  526. public static func merge<S: SignalType where S.Value == Value, S.Error == Error>(signals: S...) -> Signal<Value, Error> {
  527. return Signal.merge(signals)
  528. }
  529. }
  530. extension SignalProducerType {
  531. /// Merges the given producers into a single `SignalProducer` that will emit
  532. /// all values from each of them, and complete when all of them have
  533. /// completed.
  534. @warn_unused_result(message="Did you forget to call `start` on the producer?")
  535. public static func merge<Seq: SequenceType, S: SignalProducerType where S.Value == Value, S.Error == Error, Seq.Generator.Element == S>(producers: Seq) -> SignalProducer<Value, Error> {
  536. return SignalProducer(values: producers).flatten(.Merge)
  537. }
  538. /// Merges the given producers into a single `SignalProducer` that will emit
  539. /// all values from each of them, and complete when all of them have
  540. /// completed.
  541. @warn_unused_result(message="Did you forget to call `start` on the producer?")
  542. public static func merge<S: SignalProducerType where S.Value == Value, S.Error == Error>(producers: S...) -> SignalProducer<Value, Error> {
  543. return SignalProducer.merge(producers)
  544. }
  545. }
  546. extension SignalType where Value: SignalProducerType, Error == Value.Error {
  547. /// Returns a signal that forwards values from the latest signal sent on
  548. /// `signal`, ignoring values sent on previous inner signal.
  549. ///
  550. /// An error sent on `signal` or the latest inner signal will be sent on the
  551. /// returned signal.
  552. ///
  553. /// The returned signal completes when `signal` and the latest inner
  554. /// signal have both completed.
  555. private func switchToLatest() -> Signal<Value.Value, Error> {
  556. return Signal<Value.Value, Error> { observer in
  557. let composite = CompositeDisposable()
  558. let serial = SerialDisposable()
  559. composite += serial
  560. composite += self.observeSwitchToLatest(observer, serial)
  561. return composite
  562. }
  563. }
  564. private func observeSwitchToLatest(observer: Observer<Value.Value, Error>, _ latestInnerDisposable: SerialDisposable) -> Disposable? {
  565. let state = Atomic(LatestState<Value, Error>())
  566. return self.observe { event in
  567. switch event {
  568. case let .Next(innerProducer):
  569. innerProducer.startWithSignal { innerSignal, innerDisposable in
  570. state.modify {
  571. // When we replace the disposable below, this prevents
  572. // the generated Interrupted event from doing any work.
  573. var state = $0
  574. state.replacingInnerSignal = true
  575. return state
  576. }
  577. latestInnerDisposable.innerDisposable = innerDisposable
  578. state.modify {
  579. var state = $0
  580. state.replacingInnerSignal = false
  581. state.innerSignalComplete = false
  582. return state
  583. }
  584. innerSignal.observe { event in
  585. switch event {
  586. case .Interrupted:
  587. // If interruption occurred as a result of a new
  588. // producer arriving, we don't want to notify our
  589. // observer.
  590. let original = state.modify {
  591. var state = $0
  592. if !state.replacingInnerSignal {
  593. state.innerSignalComplete = true
  594. }
  595. return state
  596. }
  597. if !original.replacingInnerSignal && original.outerSignalComplete {
  598. observer.sendCompleted()
  599. }
  600. case .Completed:
  601. let original = state.modify {
  602. var state = $0
  603. state.innerSignalComplete = true
  604. return state
  605. }
  606. if original.outerSignalComplete {
  607. observer.sendCompleted()
  608. }
  609. case .Next, .Failed:
  610. observer.action(event)
  611. }
  612. }
  613. }
  614. case let .Failed(error):
  615. observer.sendFailed(error)
  616. case .Completed:
  617. let original = state.modify {
  618. var state = $0
  619. state.outerSignalComplete = true
  620. return state
  621. }
  622. if original.innerSignalComplete {
  623. observer.sendCompleted()
  624. }
  625. case .Interrupted:
  626. observer.sendInterrupted()
  627. }
  628. }
  629. }
  630. }
  631. extension SignalProducerType where Value: SignalProducerType, Error == Value.Error {
  632. /// Returns a signal that forwards values from the latest signal sent on
  633. /// `signal`, ignoring values sent on previous inner signal.
  634. ///
  635. /// An error sent on `signal` or the latest inner signal will be sent on the
  636. /// returned signal.
  637. ///
  638. /// The returned signal completes when `signal` and the latest inner
  639. /// signal have both completed.
  640. private func switchToLatest() -> SignalProducer<Value.Value, Error> {
  641. return SignalProducer<Value.Value, Error> { observer, disposable in
  642. let latestInnerDisposable = SerialDisposable()
  643. disposable.addDisposable(latestInnerDisposable)
  644. self.startWithSignal { signal, signalDisposable in
  645. disposable += signalDisposable
  646. disposable += signal.observeSwitchToLatest(observer, latestInnerDisposable)
  647. }
  648. }
  649. }
  650. }
  651. private struct LatestState<Value, Error: ErrorType> {
  652. var outerSignalComplete: Bool = false
  653. var innerSignalComplete: Bool = true
  654. var replacingInnerSignal: Bool = false
  655. }
  656. extension SignalType {
  657. /// Maps each event from `signal` to a new signal, then flattens the
  658. /// resulting producers (into a signal of values), according to the
  659. /// semantics of the given strategy.
  660. ///
  661. /// If `signal` or any of the created producers fail, the returned signal
  662. /// will forward that failure immediately.
  663. @warn_unused_result(message="Did you forget to call `observe` on the signal?")
  664. public func flatMap<U>(strategy: FlattenStrategy, transform: Value -> SignalProducer<U, Error>) -> Signal<U, Error> {
  665. return map(transform).flatten(strategy)
  666. }
  667. /// Maps each event from `signal` to a new signal, then flattens the
  668. /// resulting producers (into a signal of values), according to the
  669. /// semantics of the given strategy.
  670. ///
  671. /// If `signal` fails, the returned signal will forward that failure
  672. /// immediately.
  673. @warn_unused_result(message="Did you forget to call `observe` on the signal?")
  674. public func flatMap<U>(strategy: FlattenStrategy, transform: Value -> SignalProducer<U, NoError>) -> Signal<U, Error> {
  675. return map(transform).flatten(strategy)
  676. }
  677. /// Maps each event from `signal` to a new signal, then flattens the
  678. /// resulting signals (into a signal of values), according to the
  679. /// semantics of the given strategy.
  680. ///
  681. /// If `signal` or any of the created signals emit an error, the returned
  682. /// signal will forward that error immediately.
  683. @warn_unused_result(message="Did you forget to call `observe` on the signal?")
  684. public func flatMap<U>(strategy: FlattenStrategy, transform: Value -> Signal<U, Error>) -> Signal<U, Error> {
  685. return map(transform).flatten(strategy)
  686. }
  687. /// Maps each event from `signal` to a new signal, then flattens the
  688. /// resulting signals (into a signal of values), according to the
  689. /// semantics of the given strategy.
  690. ///
  691. /// If `signal` emits an error, the returned signal will forward that
  692. /// error immediately.
  693. @warn_unused_result(message="Did you forget to call `observe` on the signal?")
  694. public func flatMap<U>(strategy: FlattenStrategy, transform: Value -> Signal<U, NoError>) -> Signal<U, Error> {
  695. return map(transform).flatten(strategy)
  696. }
  697. }
  698. extension SignalType where Error == NoError {
  699. /// Maps each event from `signal` to a new signal, then flattens the
  700. /// resulting signals (into a signal of values), according to the
  701. /// semantics of the given strategy.
  702. ///
  703. /// If any of the created signals emit an error, the returned signal
  704. /// will forward that error immediately.
  705. @warn_unused_result(message="Did you forget to call `observe` on the signal?")
  706. public func flatMap<U, E>(strategy: FlattenStrategy, transform: Value -> SignalProducer<U, E>) -> Signal<U, E> {
  707. return map(transform).flatten(strategy)
  708. }
  709. /// Maps each event from `signal` to a new signal, then flattens the
  710. /// resulting signals (into a signal of values), according to the
  711. /// semantics of the given strategy.
  712. @warn_unused_result(message="Did you forget to call `observe` on the signal?")
  713. public func flatMap<U>(strategy: FlattenStrategy, transform: Value -> SignalProducer<U, NoError>) -> Signal<U, NoError> {
  714. return map(transform).flatten(strategy)
  715. }
  716. /// Maps each event from `signal` to a new signal, then flattens the
  717. /// resulting signals (into a signal of values), according to the
  718. /// semantics of the given strategy.
  719. ///
  720. /// If any of the created signals emit an error, the returned signal
  721. /// will forward that error immediately.
  722. @warn_unused_result(message="Did you forget to call `observe` on the signal?")
  723. public func flatMap<U, E>(strategy: FlattenStrategy, transform: Value -> Signal<U, E>) -> Signal<U, E> {
  724. return map(transform).flatten(strategy)
  725. }
  726. /// Maps each event from `signal` to a new signal, then flattens the
  727. /// resulting signals (into a signal of values), according to the
  728. /// semantics of the given strategy.
  729. @warn_unused_result(message="Did you forget to call `observe` on the signal?")
  730. public func flatMap<U>(strategy: FlattenStrategy, transform: Value -> Signal<U, NoError>) -> Signal<U, NoError> {
  731. return map(transform).flatten(strategy)
  732. }
  733. }
  734. extension SignalProducerType {
  735. /// Maps each event from `self` to a new producer, then flattens the
  736. /// resulting producers (into a producer of values), according to the
  737. /// semantics of the given strategy.
  738. ///
  739. /// If `self` or any of the created producers fail, the returned producer
  740. /// will forward that failure immediately.
  741. @warn_unused_result(message="Did you forget to call `start` on the producer?")
  742. public func flatMap<U>(strategy: FlattenStrategy, transform: Value -> SignalProducer<U, Error>) -> SignalProducer<U, Error> {
  743. return map(transform).flatten(strategy)
  744. }
  745. /// Maps each event from `self` to a new producer, then flattens the
  746. /// resulting producers (into a producer of values), according to the
  747. /// semantics of the given strategy.
  748. ///
  749. /// If `self` fails, the returned producer will forward that failure
  750. /// immediately.
  751. @warn_unused_result(message="Did you forget to call `start` on the producer?")
  752. public func flatMap<U>(strategy: FlattenStrategy, transform: Value -> SignalProducer<U, NoError>) -> SignalProducer<U, Error> {
  753. return map(transform).flatten(strategy)
  754. }
  755. /// Maps each event from `self` to a new producer, then flattens the
  756. /// resulting signals (into a producer of values), according to the
  757. /// semantics of the given strategy.
  758. ///
  759. /// If `self` or any of the created signals emit an error, the returned
  760. /// producer will forward that error immediately.
  761. @warn_unused_result(message="Did you forget to call `start` on the producer?")
  762. public func flatMap<U>(strategy: FlattenStrategy, transform: Value -> Signal<U, Error>) -> SignalProducer<U, Error> {
  763. return map(transform).flatten(strategy)
  764. }
  765. /// Maps each event from `self` to a new producer, then flattens the
  766. /// resulting signals (into a producer of values), according to the
  767. /// semantics of the given strategy.
  768. ///
  769. /// If `self` emits an error, the returned producer will forward that
  770. /// error immediately.
  771. @warn_unused_result(message="Did you forget to call `start` on the producer?")
  772. public func flatMap<U>(strategy: FlattenStrategy, transform: Value -> Signal<U, NoError>) -> SignalProducer<U, Error> {
  773. return map(transform).flatten(strategy)
  774. }
  775. }
  776. extension SignalProducerType where Error == NoError {
  777. /// Maps each event from `self` to a new producer, then flattens the
  778. /// resulting producers (into a producer of values), according to the
  779. /// semantics of the given strategy.
  780. ///
  781. /// If any of the created producers fail, the returned producer will
  782. /// forward that failure immediately.
  783. @warn_unused_result(message="Did you forget to call `start` on the producer?")
  784. public func flatMap<U, E>(strategy: FlattenStrategy, transform: Value -> SignalProducer<U, E>) -> SignalProducer<U, E> {
  785. return map(transform).flatten(strategy)
  786. }
  787. /// Maps each event from `self` to a new producer, then flattens the
  788. /// resulting producers (into a producer of values), according to the
  789. /// semantics of the given strategy.
  790. @warn_unused_result(message="Did you forget to call `start` on the producer?")
  791. public func flatMap<U>(strategy: FlattenStrategy, transform: Value -> SignalProducer<U, NoError>) -> SignalProducer<U, NoError> {
  792. return map(transform).flatten(strategy)
  793. }
  794. /// Maps each event from `self` to a new producer, then flattens the
  795. /// resulting signals (into a producer of values), according to the
  796. /// semantics of the given strategy.
  797. ///
  798. /// If any of the created signals emit an error, the returned
  799. /// producer will forward that error immediately.
  800. @warn_unused_result(message="Did you forget to call `start` on the producer?")
  801. public func flatMap<U, E>(strategy: FlattenStrategy, transform: Value -> Signal<U, E>) -> SignalProducer<U, E> {
  802. return map(transform).flatten(strategy)
  803. }
  804. /// Maps each event from `self` to a new producer, then flattens the
  805. /// resulting signals (into a producer of values), according to the
  806. /// semantics of the given strategy.
  807. @warn_unused_result(message="Did you forget to call `start` on the producer?")
  808. public func flatMap<U>(strategy: FlattenStrategy, transform: Value -> Signal<U, NoError>) -> SignalProducer<U, NoError> {
  809. return map(transform).flatten(strategy)
  810. }
  811. }
  812. extension SignalType {
  813. /// Catches any failure that may occur on the input signal, mapping to a new
  814. /// producer that starts in its place.
  815. @warn_unused_result(message="Did you forget to call `observe` on the signal?")
  816. public func flatMapError<F>(handler: Error -> SignalProducer<Value, F>) -> Signal<Value, F> {
  817. return Signal { observer in
  818. self.observeFlatMapError(handler, observer, SerialDisposable())
  819. }
  820. }
  821. private func observeFlatMapError<F>(handler: Error -> SignalProducer<Value, F>, _ observer: Observer<Value, F>, _ serialDisposable: SerialDisposable) -> Disposable? {
  822. return self.observe { event in
  823. switch event {
  824. case let .Next(value):
  825. observer.sendNext(value)
  826. case let .Failed(error):
  827. handler(error).startWithSignal { signal, disposable in
  828. serialDisposable.innerDisposable = disposable
  829. signal.observe(observer)
  830. }
  831. case .Completed:
  832. observer.sendCompleted()
  833. case .Interrupted:
  834. observer.sendInterrupted()
  835. }
  836. }
  837. }
  838. }
  839. extension SignalProducerType {
  840. /// Catches any failure that may occur on the input producer, mapping to a
  841. /// new producer that starts in its place.
  842. @warn_unused_result(message="Did you forget to call `start` on the producer?")
  843. public func flatMapError<F>(handler: Error -> SignalProducer<Value, F>) -> SignalProducer<Value, F> {
  844. return SignalProducer { observer, disposable in
  845. let serialDisposable = SerialDisposable()
  846. disposable.addDisposable(serialDisposable)
  847. self.startWithSignal { signal, signalDisposable in
  848. serialDisposable.innerDisposable = signalDisposable
  849. signal.observeFlatMapError(handler, observer, serialDisposable)
  850. }
  851. }
  852. }
  853. }