observableToAsyncIterable.js 5.2 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124
  1. var __assign = (this && this.__assign) || function () {
  2. __assign = Object.assign || function(t) {
  3. for (var s, i = 1, n = arguments.length; i < n; i++) {
  4. s = arguments[i];
  5. for (var p in s) if (Object.prototype.hasOwnProperty.call(s, p))
  6. t[p] = s[p];
  7. }
  8. return t;
  9. };
  10. return __assign.apply(this, arguments);
  11. };
  12. var __awaiter = (this && this.__awaiter) || function (thisArg, _arguments, P, generator) {
  13. function adopt(value) { return value instanceof P ? value : new P(function (resolve) { resolve(value); }); }
  14. return new (P || (P = Promise))(function (resolve, reject) {
  15. function fulfilled(value) { try { step(generator.next(value)); } catch (e) { reject(e); } }
  16. function rejected(value) { try { step(generator["throw"](value)); } catch (e) { reject(e); } }
  17. function step(result) { result.done ? resolve(result.value) : adopt(result.value).then(fulfilled, rejected); }
  18. step((generator = generator.apply(thisArg, _arguments || [])).next());
  19. });
  20. };
  21. var __generator = (this && this.__generator) || function (thisArg, body) {
  22. var _ = { label: 0, sent: function() { if (t[0] & 1) throw t[1]; return t[1]; }, trys: [], ops: [] }, f, y, t, g;
  23. return g = { next: verb(0), "throw": verb(1), "return": verb(2) }, typeof Symbol === "function" && (g[Symbol.iterator] = function() { return this; }), g;
  24. function verb(n) { return function (v) { return step([n, v]); }; }
  25. function step(op) {
  26. if (f) throw new TypeError("Generator is already executing.");
  27. while (_) try {
  28. if (f = 1, y && (t = op[0] & 2 ? y["return"] : op[0] ? y["throw"] || ((t = y["return"]) && t.call(y), 0) : y.next) && !(t = t.call(y, op[1])).done) return t;
  29. if (y = 0, t) op = [op[0] & 2, t.value];
  30. switch (op[0]) {
  31. case 0: case 1: t = op; break;
  32. case 4: _.label++; return { value: op[1], done: false };
  33. case 5: _.label++; y = op[1]; op = [0]; continue;
  34. case 7: op = _.ops.pop(); _.trys.pop(); continue;
  35. default:
  36. if (!(t = _.trys, t = t.length > 0 && t[t.length - 1]) && (op[0] === 6 || op[0] === 2)) { _ = 0; continue; }
  37. if (op[0] === 3 && (!t || (op[1] > t[0] && op[1] < t[3]))) { _.label = op[1]; break; }
  38. if (op[0] === 6 && _.label < t[1]) { _.label = t[1]; t = op; break; }
  39. if (t && _.label < t[2]) { _.label = t[2]; _.ops.push(op); break; }
  40. if (t[2]) _.ops.pop();
  41. _.trys.pop(); continue;
  42. }
  43. op = body.call(thisArg, _);
  44. } catch (e) { op = [6, e]; y = 0; } finally { f = t = 0; }
  45. if (op[0] & 5) throw op[1]; return { value: op[0] ? op[1] : void 0, done: true };
  46. }
  47. };
  48. Object.defineProperty(exports, "__esModule", { value: true });
  49. var iterall_1 = require("iterall");
  50. function observableToAsyncIterable(observable) {
  51. var _a;
  52. var pullQueue = [];
  53. var pushQueue = [];
  54. var listening = true;
  55. var pushValue = function (_a) {
  56. var data = _a.data;
  57. if (pullQueue.length !== 0) {
  58. pullQueue.shift()({ value: data, done: false });
  59. }
  60. else {
  61. pushQueue.push({ value: data });
  62. }
  63. };
  64. var pushError = function (error) {
  65. if (pullQueue.length !== 0) {
  66. pullQueue.shift()({ value: { errors: [error] }, done: false });
  67. }
  68. else {
  69. pushQueue.push({ value: { errors: [error] } });
  70. }
  71. };
  72. var pullValue = function () {
  73. return new Promise(function (resolve) {
  74. if (pushQueue.length !== 0) {
  75. var element = pushQueue.shift();
  76. // either {value: {errors: [...]}} or {value: ...}
  77. resolve(__assign(__assign({}, element), { done: false }));
  78. }
  79. else {
  80. pullQueue.push(resolve);
  81. }
  82. });
  83. };
  84. var subscription = observable.subscribe({
  85. next: function (value) {
  86. pushValue(value);
  87. },
  88. error: function (err) {
  89. pushError(err);
  90. },
  91. });
  92. var emptyQueue = function () {
  93. if (listening) {
  94. listening = false;
  95. subscription.unsubscribe();
  96. pullQueue.forEach(function (resolve) { return resolve({ value: undefined, done: true }); });
  97. pullQueue.length = 0;
  98. pushQueue.length = 0;
  99. }
  100. };
  101. return _a = {
  102. next: function () {
  103. return __awaiter(this, void 0, void 0, function () {
  104. return __generator(this, function (_a) {
  105. return [2 /*return*/, listening ? pullValue() : this.return()];
  106. });
  107. });
  108. },
  109. return: function () {
  110. emptyQueue();
  111. return Promise.resolve({ value: undefined, done: true });
  112. },
  113. throw: function (error) {
  114. emptyQueue();
  115. return Promise.reject(error);
  116. }
  117. },
  118. _a[iterall_1.$$asyncIterator] = function () {
  119. return this;
  120. },
  121. _a;
  122. }
  123. exports.observableToAsyncIterable = observableToAsyncIterable;
  124. //# sourceMappingURL=observableToAsyncIterable.js.map