{"_id":"channel-ts","_rev":"3-7256016418270967431924fd5776377a","name":"channel-ts","dist-tags":{"latest":"0.1.2"},"versions":{"0.1.0":{"name":"channel-ts","version":"0.1.0","description":"Channels implemented in Typescript using async/await","repository":{"type":"git","url":"git+https://github.com/arjsin/channels.git"},"main":"lib/index.js","types":"lib/index.d.ts","scripts":{"build":"tsc -p tsconfig.build.json","prepare":"tsc -p tsconfig.build.json","check":"eslint --fix  src/**/*.ts && tsc --noEmit","test":"jest"},"author":{"name":"Arjun Singh"},"license":"MIT","devDependencies":{"@types/jest":"^25.1.2","@types/node":"^13.7.1","@typescript-eslint/eslint-plugin":"^2.19.2","@typescript-eslint/parser":"^2.19.2","eslint":"^6.8.0","jest":"^25.1.0","reflect-metadata":"^0.1.13","ts-jest":"^25.2.0","typescript":"^3.7.5"},"jest":{"preset":"ts-jest","testEnvironment":"node"},"gitHead":"1967ed3f208196c1d3b23727e1e44f9d431c7319","bugs":{"url":"https://github.com/arjsin/channels/issues"},"homepage":"https://github.com/arjsin/channels#readme","_id":"channel-ts@0.1.0","_nodeVersion":"13.9.0","_npmVersion":"6.14.2","dist":{"integrity":"sha512-9RTC9JjCf/7aaeReBNTnqtNFOjUnDqN9HifQxoXF/7zQjy4jyhzh8Z3UDGkeVC8ZbECobWbo4CJKon8NQIahAw==","shasum":"9f4dbf2789e49318c7586dd4f9b14a1731aac378","tarball":"https://registry.npmjs.org/channel-ts/-/channel-ts-0.1.0.tgz","fileCount":9,"unpackedSize":13648,"npm-signature":"-----BEGIN PGP SIGNATURE-----\r\nVersion: OpenPGP.js v3.0.4\r\nComment: https://openpgpjs.org\r\n\r\nwsFcBAEBCAAQBQJeltLSCRA9TVsSAnZWagAAjzkP/1kNF5Am99ZaQVOra+/l\nbOiqL0DnuJ1MOc1NCDXEUyRD1bexn0C/odUEa8CyuODWPp8lY7r4ZmFmeOpK\n8ke6//8ztmhqUZ4iNx5XapqFeG+/heD6GjLmRPH/w8uftXMexeLVbpB/FoXH\nWd0vtsFgRXwTfqpJivq/0auxQMU8R+WJkFHY7d1PBiYBhgyeSdfWyPXvQYYH\nWtwzAmCgDFwHrlTvvH1kAor+CotmonNk+XmtyM2R1jvZ4x4kyFbUAytgWppz\n0qMW7lxVgqXTZ1VQFHByLoDZ8JJcnGkF3FptOwFyffdjhzs4wW3bKat+f7t6\njzCrGFVjnhcRG+5imuQLzuKVd+RFyQRhGRUuIagnORr8Mz3h/NP/Hw2az2E0\n8ZIn0dnWBofCBDALLlPiN0G3dTiDyLDtrIOX6KMzRLD7/818tEJqceTcV0ei\nMdw/T0eg/jzabH1SSDliH/akBxmzKaMc3d/2ghKovyi1LjTOQuA3K0SbOnEF\nbDmb+v5cF5nvFLtPfoBT36YIw13jjvS8Qwiohp532si4+b9YCfIzlRRll+fY\n8qSOKn4DUUzaSy/tBDeXDeHG6lqqp90+r6VPcxA+hH7Ge2+Ee08Sb8kOF8ca\ndMApMhvpcBeZsk9Yzx+iLqvaVm3EybkiQdyw6YT99pdbI7oYpxCc1CNTb+qt\ne4/s\r\n=ORz7\r\n-----END PGP SIGNATURE-----\r\n","signatures":[{"keyid":"SHA256:jl3bwswu80PjjokCgh0o2w5c2U4LhQAE57gj9cz1kzA","sig":"MEUCIDdpiuVRN73VqjKBnIw28S7tT+YV1WBpW3XPkKYkYvQZAiEAxN5qeEtS/zRhU4t/7edrHU9JjjuGIFIUW+3qos6AEw8="}]},"maintainers":[{"name":"arjsin","email":"arj1singh@gmail.com"}],"_npmUser":{"name":"arjsin","email":"arj1singh@gmail.com"},"directories":{},"_npmOperationalInternal":{"host":"s3://npm-registry-packages","tmp":"tmp/channel-ts_0.1.0_1586942673563_0.774189349390525"},"_hasShrinkwrap":false},"0.1.1":{"name":"channel-ts","version":"0.1.1","description":"Channels implemented in Typescript using async/await","repository":{"type":"git","url":"git+https://github.com/arjsin/channels.git"},"main":"lib/index.js","types":"lib/index.d.ts","scripts":{"build":"tsc -p tsconfig.build.json","prepare":"tsc -p tsconfig.build.json","check":"eslint --fix  src/**/*.ts && tsc --noEmit","test":"jest"},"author":{"name":"Arjun Singh"},"license":"MIT","devDependencies":{"@types/jest":"^25.1.2","@types/node":"^13.7.1","@typescript-eslint/eslint-plugin":"^2.19.2","@typescript-eslint/parser":"^2.19.2","eslint":"^6.8.0","jest":"^25.1.0","reflect-metadata":"^0.1.13","ts-jest":"^25.2.0","typescript":"^3.7.5"},"keywords":["channel","csp","typescript","async","await","observer","observable","publisher","subscriber"],"jest":{"preset":"ts-jest","testEnvironment":"node"},"gitHead":"668a4e8ec752c1f8881205df219a2ce99ee0001c","bugs":{"url":"https://github.com/arjsin/channels/issues"},"homepage":"https://github.com/arjsin/channels#readme","_id":"channel-ts@0.1.1","_nodeVersion":"13.9.0","_npmVersion":"6.14.2","dist":{"integrity":"sha512-AKUsmcHwokBzdSgxnpm9xT2FPuWuV5nCjjNI7weOCczpi5PbvmzSKuA90dIQFVJFpvHEycHZ5CR5wRzPQ0iPuA==","shasum":"6c5c27cdc4ce0f4fdb1ff43da207f1f71fab857d","tarball":"https://registry.npmjs.org/channel-ts/-/channel-ts-0.1.1.tgz","fileCount":9,"unpackedSize":15143,"npm-signature":"-----BEGIN PGP SIGNATURE-----\r\nVersion: OpenPGP.js v3.0.4\r\nComment: https://openpgpjs.org\r\n\r\nwsFcBAEBCAAQBQJepWTwCRA9TVsSAnZWagAAkz4QAJsQfL7MsksRvHwhAlIb\n6HOp8kTYR5EjwDt/mezh0DShjbbXXh2h0QhvPoYKw0mGPzIvJT5Km1SZ0FdU\nbRkZaPTt8is29EJMrNbKEiktfaWJjmUdwc8lVvtzh/ShdfoIASpkSMAFNQYg\n6LMy8vlQx3GwL9Pe7Ag6iaBgNXe1wAa7gquUA+yWw6dQt7NCSPOmdLtju+Ng\nDCRnbn0/9UbmCax/mB7Qowc02/P9olnUry+S73hUuWeb8t5s5ZzipzhRxUkw\ndvWDjwkK2PkTDlHlZkKP+/NNqq/WM11C0EsdklJ6Dq+0Jl7bfo7oMKgj1txI\naKrE35NzORnuxzQS9PwVhtvprsC0crsnKpJ3Z1sGbmFyZ5Brxi9go7llob08\neq/0dUiDtHjcBCTucwBtRRFuRqEHqmdy5dHXHUK2oc54HRlMuCNTF8+sQTGA\nCOy5aJ4986E49ZaVGngQx73aWrFo8nFxdvVngJM7WwpYO18kSgUSIuc/UvWW\ny/UIL4KYMPuNr+ocRFUpPd8dN2S935OeBMnxtYsKHl8cwXHpzApVN2IrHNbF\ndBskw2OKL9AWovJZDrg1j5JbpouWqk26djIwGN4qemCwb/h8FF9CBQmbVbyz\n30SEeu324KcoD9PYbj8EJsLjkyWS9tVQOl0zQGqyC+LezjtSbAa2GQQ7RtV+\nzyB8\r\n=76Sh\r\n-----END PGP SIGNATURE-----\r\n","signatures":[{"keyid":"SHA256:jl3bwswu80PjjokCgh0o2w5c2U4LhQAE57gj9cz1kzA","sig":"MEYCIQD3/2SnFanWNYtZLpBDioOLbNFvGhAIvYkmjDmfyS06FwIhANMXENJpa2S+51vsCdRcuDkFixIbOjWe9XlLd9bYR5fy"}]},"maintainers":[{"name":"arjsin","email":"arj1singh@gmail.com"}],"_npmUser":{"name":"arjsin","email":"arj1singh@gmail.com"},"directories":{},"_npmOperationalInternal":{"host":"s3://npm-registry-packages","tmp":"tmp/channel-ts_0.1.1_1587897583962_0.14286393429404987"},"_hasShrinkwrap":false},"0.1.2":{"name":"channel-ts","version":"0.1.2","description":"Channels implemented in Typescript using async/await","repository":{"type":"git","url":"git+https://github.com/arjsin/channels.git"},"main":"lib/index.js","types":"lib/index.d.ts","scripts":{"build":"tsc -p tsconfig.build.json","prepare":"tsc -p tsconfig.build.json","check":"eslint --fix  src/**/*.ts && tsc --noEmit","test":"jest"},"author":{"name":"Arjun Singh"},"license":"MIT","devDependencies":{"@types/jest":"^25.1.2","@types/node":"^13.7.1","@typescript-eslint/eslint-plugin":"^2.19.2","@typescript-eslint/parser":"^2.19.2","eslint":"^6.8.0","jest":"^25.1.0","reflect-metadata":"^0.1.13","ts-jest":"^25.2.0","typescript":"^3.9.7"},"keywords":["channel","csp","typescript","async","await","observer","observable","publisher","subscriber","mutex"],"jest":{"preset":"ts-jest","testEnvironment":"node"},"gitHead":"0842037044c5265ee0536d626674860d243b66bc","bugs":{"url":"https://github.com/arjsin/channels/issues"},"homepage":"https://github.com/arjsin/channels#readme","_id":"channel-ts@0.1.2","_nodeVersion":"14.4.0","_npmVersion":"6.14.5","dist":{"integrity":"sha512-cI/XiDF+jB0v95Xup8xlM7k93lT3xwPl0WdjEZ9w9aUMf5N+3GQevspK2EDYfMyxcKcXdN1F6PDpuYRpUfaZmg==","shasum":"eefdbc7c9e0275aa1862f8c436cca1c49658f768","tarball":"https://registry.npmjs.org/channel-ts/-/channel-ts-0.1.2.tgz","fileCount":13,"unpackedSize":20028,"npm-signature":"-----BEGIN PGP SIGNATURE-----\r\nVersion: OpenPGP.js v3.0.4\r\nComment: https://openpgpjs.org\r\n\r\nwsFcBAEBCAAQBQJfOiXeCRA9TVsSAnZWagAAMisP/2IY2akMN9ny35EyekUR\n/mUuiIGW7+fsf2eahxz/gyhOQMQ/Ltid13qgqoD3bZeHLn8c8ysZ7/tp+9WL\nPUn75gbb56YPoJ5a7UR4Ck2IyOV7Sb9f7TBVld49Fgh5U58JJHbO81S1sFPx\nORQYxCy+AmDo0EIEDzTk41UoJ3LMsa56BUPxqSOGh5/SkD4HcXCw9JXD3Cq/\n6ixmGqrq7a54LKfoM4VCixWwAd8ZAEL0R6d0iFiA3sbMBWKcdM0pbghaNMha\n4UM5sS0DrKKzigW7LXsHV0G/ywCeV/aZwqt9Sbn0SBeezKnykBvUlEbQ+sSD\n4XejwWeGXTQgOGWts+PArju/mzvmFtoAECVTrZEaFMgoVVsGSR/K5Vwu8oGS\n6VMfUbigbhbalCUbCPXXyPVATh5S5dWl6elN57UA62kNJGupnaD0BP+Hxyn9\nUwHiPxvE3uN4lkhjGXTRmSiczczogvdG4c2x1jCbNQlQiUwtM2ZIqJ45/q4L\nFYqrbF2Xh0n80yHRNeEUEIZsC7JKLHblvIoRbIaEhVz98sU2qkj+LOgHJidY\n0yMgr7vTJE8n5RC8KhXYo/zAseY0K9jcuO/1+lgNF8h47VTSPmJQnoWqhkDL\nCOb/7ZDokgS4XkAxYm0HIvpTlgDVdNtdSTlkcCRQx8FjZG1POVIRiObmuyeH\nj7EF\r\n=0wRb\r\n-----END PGP SIGNATURE-----\r\n","signatures":[{"keyid":"SHA256:jl3bwswu80PjjokCgh0o2w5c2U4LhQAE57gj9cz1kzA","sig":"MEQCIC0USKbZFscUOZJYIQEd/XFDCMsiXb+Vob0qGpt2I0N7AiBp+piIiBvYTjy0Xqw4/frnDuBq2lnk7mHOANQbXhta4A=="}]},"maintainers":[{"name":"arjsin","email":"arj1singh@gmail.com"}],"_npmUser":{"name":"arjsin","email":"arj1singh@gmail.com"},"directories":{},"_npmOperationalInternal":{"host":"s3://npm-registry-packages","tmp":"tmp/channel-ts_0.1.2_1597646301636_0.3957179905111532"},"_hasShrinkwrap":false}},"time":{"created":"2020-04-15T09:24:33.563Z","0.1.0":"2020-04-15T09:24:33.696Z","modified":"2022-04-12T05:55:44.341Z","0.1.1":"2020-04-26T10:39:44.099Z","0.1.2":"2020-08-17T06:38:21.731Z"},"maintainers":[{"name":"arjsin","email":"arj1singh@gmail.com"}],"description":"Channels implemented in Typescript using async/await","homepage":"https://github.com/arjsin/channels#readme","repository":{"type":"git","url":"git+https://github.com/arjsin/channels.git"},"author":{"name":"Arjun Singh"},"bugs":{"url":"https://github.com/arjsin/channels/issues"},"license":"MIT","readme":"# channel-ts\n>Minimal Async/Await Channels in Typescript\n\n[![Node.js CI](https://img.shields.io/github/workflow/status/arjsin/channels/Node.js%20CI?style=flat-square)](https://github.com/arjsin/channels/actions?query=workflow%3A%22Node.js+CI%22)\n[![License](https://img.shields.io/:license-mit-blue.svg?style=flat-square)](/LICENSE)\n[![NPM](https://img.shields.io/npm/v/channel-ts?style=flat-square)](https://www.npmjs.com/package/channel-ts)\n\n## Features\n- Simple API with JavaScript's async iterators for receivers\n- Broadcast on MultiReceiverChannel\n- Observe on object and notify manually\n- Mutex\n\n## Examples\n### Multi producer and single consumer\nThis channel allows multiple sender to send data to a single receiver.\nThe messages starts buffering as soon as the channel is created.\nNo messages are lost.\n\n```typescript\nimport { SimpleChannel } from \"channel-ts\";\n\nconst delay = (ms: number) => new Promise((resolve) => setTimeout(resolve, ms));\n\n// printer waits for the messages on the channel until it closes\nasync function printer(chan: SimpleChannel<string>) {\n    for await(const data of chan) { // use async iterator to receive data\n        console.log(`Received: ${data}`);\n    }\n    console.log(\"Closed\");\n}\n\n// sender sends some messages to the channel\nasync function sender(id: number, chan: SimpleChannel<string>) {\n    await delay(id*2000);\n    chan.send(`hello from ${id}`); // sends data, boundless channels don't block\n    await delay(2800);\n    chan.send(`bye from ${id}`); // sends some data again\n}\n\nasync function main() {\n    const chan = new SimpleChannel<string>(); // creates a new simple channel\n    const p1 = printer(chan); // uses the channel to print the received data\n    const p2 = [0, 1, 2, 3, 4].map(async i => sender(i, chan)); // creates and spawns senders\n\n    await Promise.all(p2); // waits for the sender\n    chan.close(); // closes the channel on the server end\n    await p1; // waits for the channel to close on the receiver end too\n}\n\nmain();\n```\n\n### Output\n[![Simple Output](../assets/simple_output.svg?raw=true&sanitize=true)](#)\n\n## Multi producer and multi consumer\nThis channel allows multiple senders to send data to a multiple receivers.\nThis channel needs explicit creation of receiver.\nAll the messages are broadcast and buffered for receivers to receive.\nSo messages are lost for the duration when the receiver was not created.\nCreation of receiver is similar to subscription in publisher-subscriber pattern.\n\n```typescript\nimport { MultiReceiverChannel, SimpleReceiver } from \"channel-ts\";\n\nconst delay = (ms: number) => new Promise((resolve) => setTimeout(resolve, ms));\n\n// printer waits for the messages on the channel until it closes\nasync function printer(id: string, chan: SimpleReceiver<string>) {\n    for await(const data of chan) { // use async iterator to receive data\n        console.log(`Printer ${id} received: ${data}`);\n    }\n    console.log(`Printer ${id} closed`);\n}\n\n// sender sends some messages to the channel\nasync function sender(id: number, chan: MultiReceiverChannel<string>) {\n    await delay(id*2000);\n    chan.send(`hello from ${id}`); // sends data, boundless channels don't block\n    await delay(2800);\n    chan.send(`bye from ${id}`); // sends some data again\n}\n\nasync function main() {\n    const chan = new MultiReceiverChannel<string>(); // creates a new simple channel\n    const r1 = chan.receiver();\n    const p1 = printer(\"A\", r1); // uses the channel to print the received data\n    const r2 = chan.receiver();\n    const p2 = printer(\"B\", r2); // uses the channel to print the received data\n    const p3 = [0, 1, 2, 3, 4].map(async i => sender(i, chan)); // create and spawn senders\n\n    await Promise.all(p3); // wait for sender\n    chan.removeReceiver(r1); // close channel\n    chan.removeReceiver(r2); // close channel\n    chan.close();\n    await Promise.all([p1, p2]); // wait for channel to close on receiver end\n}\n\nmain();\n```\n\n### Output\n[![Simple Output](../assets/multi_output.svg?raw=true&sanitize=true)](#)\n\n## Observe\nObserve is a function which creates Observable type of JavaScript objects.\nThe object is mutated and `notify()` is called to send data to receiver.\nThis technique is similar to observer pattern.\n\nPlease note that the object is shared between the sender and the receiver, so the receiver while reading the received message might run into a risk of observing further changes made by the sender.\nTo guarantee that a received message is not modified, please make sure that each receiver is not pre-empted for the entire duration of reading the message and the message is not mutated by the receivers.\n\n```typescript\nimport { observe, SimpleReceiver, Observable } from \"channel-ts\";\n\nconst delay = (ms: number) => new Promise((resolve) => setTimeout(resolve, ms));\n\n// printer waits for the messages on the channel until it closes\nasync function printer(id: string, chan: SimpleReceiver<number[]>) {\n    for await(const data of chan) { // use async iterator to receive data\n        console.log(`Printer ${id} received: ${JSON.stringify(data)}`);\n    }\n    console.log(`Printer ${id} closed`);\n}\n\n// sender sends some messages to the channel\nasync function sender(id: number, array: Observable<number[]>) {\n    await delay(id*1500);\n    array.fill(1); // does some manipulation\n    array.notify(); // notifies all the receivers with the value\n    await delay(2200);\n    array[0] = id * 111; // does some manipulation\n    array[1] = 0;\n    array[2] = (9-id) * 111;\n    array.notify(); // notifies all the receivers with the value\n}\n\nasync function main() {\n    const chan = observe([0, 0, 0]); // creates a new observable, works with objects\n    const r1 = chan.receiver();\n    const p1 = printer(\"A\", r1); // uses the channel to print received data\n    const r2 = chan.receiver();\n    const p2 = printer(\"B\", r2); // uses the channel to print received data\n    const p3 = [0, 1, 2, 3, 4].map(async i => sender(i, chan)); // create and spawn senders\n\n    await Promise.all(p3); // wait for sender\n    chan.removeReceiver(r1); // close channel\n    chan.removeReceiver(r2); // close channel\n    await Promise.all([p1, p2]); // wait for channel to close on receiver end\n}\n\nmain();\n```\n### Output\n[![Simple Output](../assets/observe_output.svg?raw=true&sanitize=true)](#)\n\n## Mutex\nMutex allows asynchronous program to have synchronization by the application of locking. This is useful when a shared resource is accessed concurrently.\n```typescript\nimport { Mutex } from \"channel-ts\";\n\nconst delay = (ms: number) => new Promise((resolve) => setTimeout(resolve, ms));\n\nconst main = async () => {\n\tconst task1 = async (): Promise<void> => {\n\t\tconsole.log(\"task 1: going to start\");\n\t\tconst guard = await mutex.acquire();\n\t\tconsole.log(\"task 1: mutex acquired\");\n\t\tawait delay(2000);  // assume shared access done by task 1\n\t\tconsole.log(\"task 1: finished\");\n\t\tguard.release();\n\t\tconsole.log(\"task 1: mutex released\");\n\t};\n\tconst task2 = async (): Promise<void> => {\n\t\tawait delay(200);\n\t\tconsole.log(\"task 2: going to start\");\n\t\tconst guard = await mutex.acquire();\n\t\tconsole.log(\"task 2: mutex acquired\");\n\t\tawait delay(1500);  // assume shared access done by task 2\n\t\tconsole.log(\"task 2: finished\");\n\t\tguard.release();\n\t\tconsole.log(\"task 2: mutex released\");\n\t};\n\tconst task3 = async (): Promise<void> => {\n\t\tawait delay(400);\n\t\tconsole.log(\"task 3: going to start\");\n\t\tconst guard = await mutex.acquire();\n\t\tconsole.log(\"task 3: mutex acquired\");\n\t\tawait delay(1800);  // assume shared access done by task 3\n\t\tconsole.log(\"task 3: finished\");\n\t\tguard.release();\n\t\tconsole.log(\"task 3: mutex released\");\n\t};\n\t// wait for all our tasks\n\tawait Promise.all([task1(), task2(), task3()]);\n};\n\nmain();\n```\n### Output\n[![Simple Output](../assets/mutex_output.svg?raw=true&sanitize=true)](#)\n\n## License\n **[MIT license](/LICENSE)**\n","readmeFilename":"README.md","keywords":["channel","csp","typescript","async","await","observer","observable","publisher","subscriber","mutex"]}