{"_id":"@dafaz/change_stream_broker","_rev":"12-6b3709dbc10375c05b46fadc77684164","name":"@dafaz/change_stream_broker","dist-tags":{"latest":"1.1.0"},"versions":{"1.0.0":{"name":"@dafaz/change_stream_broker","version":"1.0.0","keywords":["mongodb","change-stream","resume-token","event-driven","microservices"],"author":{"name":"asammarco"},"license":"MIT","_id":"@dafaz/change_stream_broker@1.0.0","maintainers":[{"name":"asammarco","email":"asammarco@hotmail.com"}],"dist":{"shasum":"2c458415778f5cd3f0414f0f406b922813632f66","tarball":"https://registry.npmjs.org/@dafaz/change_stream_broker/-/change_stream_broker-1.0.0.tgz","fileCount":46,"integrity":"sha512-0L1OosiM4oh6NNeIYLsyiWwED5bzxKlmlTgc762SajKG0cmx9V6mMxinrtTY+EUf3rFsqKoQpTtQJ/wrJ39OdQ==","signatures":[{"sig":"MEQCICLB2xhYd7ajKXAFWxUTySdFbL/i7mAWZHV6r3/2bbLDAiAZMElBKjcUQXMX4pW7G8eGvtHej7xVi31Hgiz9q0w2Xw==","keyid":"SHA256:DhQ8wR5APBvFHLF/+Tc+AYvPOdTpcIDqOhxsBHRwC7U"}],"unpackedSize":120487},"main":"dist/index.js","types":"dist/index.d.ts","gitHead":"904219b7e41dd412911bb3dab37e8587ebcb83be","private":false,"scripts":{"lint":"npx @biomejs/biome check --write","test":"echo \"Error: no test specified\" && exit 1","build":"tsc"},"_npmUser":{"name":"asammarco","email":"asammarco@hotmail.com"},"_npmVersion":"10.9.0","description":"A reusable MongoDB Change Stream manager","directories":{},"_nodeVersion":"23.2.0","dependencies":{"mongodb":"^6.19.0","typescript":"^5.9.2"},"_hasShrinkwrap":false,"devDependencies":{"@types/node":"^24.3.0","@biomejs/biome":"2.2.2"},"peerDependencies":{"mongodb":"^6.19.0"},"_npmOperationalInternal":{"tmp":"tmp/change_stream_broker_1.0.0_1756493239627_0.9921386022170327","host":"s3://npm-registry-packages-npm-production"},"deprecated":"Package no longer supported. Contact Support at https://www.npmjs.com/support for more info."},"1.0.1":{"name":"@dafaz/change_stream_broker","version":"1.0.1","keywords":["mongodb","change-stream","resume-token","event-driven","microservices"],"author":{"name":"asammarco"},"license":"MIT","_id":"@dafaz/change_stream_broker@1.0.1","maintainers":[{"name":"asammarco","email":"asammarco@hotmail.com"}],"dist":{"shasum":"cb7219630c42a149df17292d5f4d52aa72fd55dd","tarball":"https://registry.npmjs.org/@dafaz/change_stream_broker/-/change_stream_broker-1.0.1.tgz","fileCount":46,"integrity":"sha512-ixSFVHXHqajzIKs6sPgI/qaTzJwHWQ7TCkih8G0o1Sj37m1A7RWHAfSkuHXA9FPm2E2B1DMP83zKChl3OLk9vg==","signatures":[{"sig":"MEQCIDDyqEPdaHrzUMkzCauk0WijLC85lvYgPDj3cehRVYdWAiARM20JuSpGvhRAV5URlxgEDOAgBEDdZd81ztuIt+dS1w==","keyid":"SHA256:DhQ8wR5APBvFHLF/+Tc+AYvPOdTpcIDqOhxsBHRwC7U"}],"unpackedSize":120487},"main":"dist/index.js","types":"dist/index.d.ts","gitHead":"f29738be017460bf6d428fa5bd96c6a590c38550","private":false,"scripts":{"lint":"npx @biomejs/biome check --write","test":"echo \"Error: no test specified\" && exit 1","build":"tsc"},"_npmUser":{"name":"asammarco","email":"asammarco@hotmail.com"},"_npmVersion":"10.9.0","description":"A reusable MongoDB Change Stream manager","directories":{},"_nodeVersion":"23.2.0","dependencies":{"mongodb":"^6.19.0","typescript":"^5.9.2"},"_hasShrinkwrap":false,"devDependencies":{"@types/node":"^24.3.0","@biomejs/biome":"2.2.2"},"peerDependencies":{"mongodb":"^6.19.0"},"_npmOperationalInternal":{"tmp":"tmp/change_stream_broker_1.0.1_1756500323969_0.5993937813265175","host":"s3://npm-registry-packages-npm-production"},"deprecated":"Package no longer supported. Contact Support at https://www.npmjs.com/support for more info."},"1.0.2":{"name":"@dafaz/change_stream_broker","version":"1.0.2","keywords":["mongodb","change-stream","resume-token","event-driven","microservices"],"author":{"name":"asammarco"},"license":"MIT","_id":"@dafaz/change_stream_broker@1.0.2","maintainers":[{"name":"asammarco","email":"asammarco@hotmail.com"}],"dist":{"shasum":"3a797931b255433ea7b4543a40bb20d9b271b0e4","tarball":"https://registry.npmjs.org/@dafaz/change_stream_broker/-/change_stream_broker-1.0.2.tgz","fileCount":46,"integrity":"sha512-ydTHWH+q3PNKdQXQp2J4xnf60r8Ub/kWo6W3H62NT+kJbmtpUF82GSl0To5M0G3uMaAT6MYeZAOcWLlObFNCcQ==","signatures":[{"sig":"MEQCID4weX69CYywbJMknhFuPC2lp/WPkshB+NArMjpcqkhLAiAOzrW6d+pKRvleUdXaSnDT/0+GYqd8cee+/EPcpc+V4w==","keyid":"SHA256:DhQ8wR5APBvFHLF/+Tc+AYvPOdTpcIDqOhxsBHRwC7U"}],"unpackedSize":120487},"main":"dist/index.js","types":"dist/index.d.ts","gitHead":"25869bac544d0689d33e92d3ca52bf955979f524","private":false,"scripts":{"lint":"npx @biomejs/biome check --write","test":"echo \"Error: no test specified\" && exit 1","build":"tsc"},"_npmUser":{"name":"asammarco","email":"asammarco@hotmail.com"},"_npmVersion":"10.9.0","description":"A reusable MongoDB Change Stream manager","directories":{},"_nodeVersion":"23.2.0","dependencies":{"mongodb":"^6.19.0","typescript":"^5.9.2"},"_hasShrinkwrap":false,"devDependencies":{"@types/node":"^24.3.0","@biomejs/biome":"2.2.2"},"peerDependencies":{"mongodb":"^6.19.0"},"_npmOperationalInternal":{"tmp":"tmp/change_stream_broker_1.0.2_1756500905154_0.7604209775269626","host":"s3://npm-registry-packages-npm-production"},"deprecated":"Package no longer supported. Contact Support at https://www.npmjs.com/support for more info."},"1.0.3":{"name":"@dafaz/change_stream_broker","version":"1.0.3","keywords":["mongodb","change-stream","resume-token","event-driven","microservices"],"author":{"name":"asammarco"},"license":"MIT","_id":"@dafaz/change_stream_broker@1.0.3","maintainers":[{"name":"asammarco","email":"asammarco@hotmail.com"}],"dist":{"shasum":"a2a90eae22bf3eb71fc5db29e5fec73553029a53","tarball":"https://registry.npmjs.org/@dafaz/change_stream_broker/-/change_stream_broker-1.0.3.tgz","fileCount":50,"integrity":"sha512-xhlZbSVLP2LQqizzpysApRDUEpTf1zpAAAIcq+UvzmwHTTDQo/w7A6mxHz1PjImUDHbkxmeKLxvvIsAk8rn0HA==","signatures":[{"sig":"MEUCIQD2YcWCsxGDQTIdB1WFeouKevzn5bdYw/f07WLxBVNfpAIgP70u53wRmD0jQSo1FpMwStGzG6jVctFq0zjSAJY3XYU=","keyid":"SHA256:DhQ8wR5APBvFHLF/+Tc+AYvPOdTpcIDqOhxsBHRwC7U"}],"unpackedSize":122698},"main":"dist/index.js","types":"dist/index.d.ts","gitHead":"25869bac544d0689d33e92d3ca52bf955979f524","private":false,"scripts":{"lint":"npx @biomejs/biome check --write","test":"echo \"Error: no test specified\" && exit 1","build":"tsc"},"_npmUser":{"name":"asammarco","email":"asammarco@hotmail.com"},"_npmVersion":"10.9.0","description":"A reusable MongoDB Change Stream manager","directories":{},"_nodeVersion":"23.2.0","dependencies":{"mongodb":"^6.19.0","typescript":"^5.9.2"},"_hasShrinkwrap":false,"devDependencies":{"@types/node":"^24.3.0","@biomejs/biome":"2.2.2"},"peerDependencies":{"mongodb":"^6.19.0"},"_npmOperationalInternal":{"tmp":"tmp/change_stream_broker_1.0.3_1756501324641_0.44516719405704896","host":"s3://npm-registry-packages-npm-production"},"deprecated":"Package no longer supported. Contact Support at https://www.npmjs.com/support for more info."},"1.0.4":{"name":"@dafaz/change_stream_broker","version":"1.0.4","keywords":["mongodb","change-stream","resume-token","event-driven","microservices"],"author":{"name":"asammarco"},"license":"MIT","_id":"@dafaz/change_stream_broker@1.0.4","maintainers":[{"name":"asammarco","email":"asammarco@hotmail.com"}],"dist":{"shasum":"b453b5c08f3276ada08c73125b724a6a566c1e71","tarball":"https://registry.npmjs.org/@dafaz/change_stream_broker/-/change_stream_broker-1.0.4.tgz","fileCount":50,"integrity":"sha512-Hfv1I3OSC+uugRVr8X1N9Pfw5xkV14qIDQPPVin0ya1zV1e/1e7jfXeiAie8kNvPAvGJC9eQ8gGDBFfScxCTDg==","signatures":[{"sig":"MEYCIQDgxOT18qv+n2f2NEukWnq0YEq80eI4FqLlOig74937HgIhAP19ca7PGi8hfRsjffoM2U2V1dLXUr/Z943DTgQfe6/2","keyid":"SHA256:DhQ8wR5APBvFHLF/+Tc+AYvPOdTpcIDqOhxsBHRwC7U"}],"unpackedSize":122740},"main":"dist/index.js","types":"dist/index.d.ts","gitHead":"fdd1acca585ae06699a6264577503a2051274c31","private":false,"scripts":{"lint":"npx @biomejs/biome check --write","test":"echo \"Error: no test specified\" && exit 1","build":"tsc && tsc-alias"},"_npmUser":{"name":"asammarco","email":"asammarco@hotmail.com"},"_npmVersion":"10.9.0","description":"A reusable MongoDB Change Stream manager","directories":{},"_nodeVersion":"23.2.0","dependencies":{"mongodb":"^6.19.0","typescript":"^5.9.2"},"_hasShrinkwrap":false,"devDependencies":{"tsc-alias":"^1.8.16","@types/node":"^24.3.0","@biomejs/biome":"2.2.2"},"peerDependencies":{"mongodb":"^6.19.0"},"_npmOperationalInternal":{"tmp":"tmp/change_stream_broker_1.0.4_1756502211346_0.9539221854344395","host":"s3://npm-registry-packages-npm-production"},"deprecated":"Package no longer supported. Contact Support at https://www.npmjs.com/support for more info."},"1.0.5":{"name":"@dafaz/change_stream_broker","version":"1.0.5","keywords":["mongodb","change-stream","resume-token","event-driven","microservices"],"author":{"name":"asammarco"},"license":"MIT","_id":"@dafaz/change_stream_broker@1.0.5","maintainers":[{"name":"asammarco","email":"asammarco@hotmail.com"}],"dist":{"shasum":"2d748a41276a537943fdd2be7e7bcdd099443653","tarball":"https://registry.npmjs.org/@dafaz/change_stream_broker/-/change_stream_broker-1.0.5.tgz","fileCount":50,"integrity":"sha512-WWcg5PWA6dcrRzTwY/IC80cbICPssJTl8VhI7nxTwTzD0anlpJRLe9gi0Mbzi/Rt8TBcv7zCyHuzIMHUrP7cAg==","signatures":[{"sig":"MEUCIBNAB9zz6cAlGaVI7DSrDD/1Lrb/rW2ml2C9Lteg4wx8AiEA1HPNkHMNH1p24Ydyp/hiXyc6fZkcQe3lYUM6eWKSIY0=","keyid":"SHA256:DhQ8wR5APBvFHLF/+Tc+AYvPOdTpcIDqOhxsBHRwC7U"}],"unpackedSize":122740},"main":"dist/index.js","types":"dist/index.d.ts","gitHead":"28e5f89bd46bbb049af0a3a533ecd24b392df070","private":false,"scripts":{"lint":"npx @biomejs/biome check --write","test":"echo \"Error: no test specified\" && exit 1","build":"tsc && tsc-alias"},"_npmUser":{"name":"asammarco","email":"asammarco@hotmail.com"},"_npmVersion":"10.9.0","description":"A reusable MongoDB Change Stream manager","directories":{},"_nodeVersion":"23.2.0","dependencies":{"mongodb":"^6.19.0","typescript":"^5.9.2"},"_hasShrinkwrap":false,"devDependencies":{"tsc-alias":"^1.8.16","@types/node":"^24.3.0","@biomejs/biome":"2.2.2"},"peerDependencies":{"mongodb":"^6.19.0"},"_npmOperationalInternal":{"tmp":"tmp/change_stream_broker_1.0.5_1756508715260_0.1575792712533166","host":"s3://npm-registry-packages-npm-production"},"deprecated":"Package no longer supported. Contact Support at https://www.npmjs.com/support for more info."},"1.0.6":{"name":"@dafaz/change_stream_broker","version":"1.0.6","keywords":["mongodb","change-stream","resume-token","event-driven","microservices"],"author":{"name":"asammarco"},"license":"MIT","_id":"@dafaz/change_stream_broker@1.0.6","maintainers":[{"name":"asammarco","email":"asammarco@hotmail.com"}],"dist":{"shasum":"605b07483826fc27e7138ef5c472218a02222d35","tarball":"https://registry.npmjs.org/@dafaz/change_stream_broker/-/change_stream_broker-1.0.6.tgz","fileCount":50,"integrity":"sha512-SHCGafq2Md8z4uKbHA7fjVpN1T9/6j8yt4BwY+SObQXTovzfjXR85Ega8vvIOhoWKCYY96qaaQ35fL6EVOJuZw==","signatures":[{"sig":"MEUCICTbdidp51yUokL7nSFDdFnQzFCP+weVABvj7CyYremsAiEA+DAM2TgZgImaaHSryUugS07+KVhjPYPEuJlHImcR8E0=","keyid":"SHA256:DhQ8wR5APBvFHLF/+Tc+AYvPOdTpcIDqOhxsBHRwC7U"}],"unpackedSize":124289},"main":"dist/index.js","types":"dist/index.d.ts","gitHead":"bc3405ca5148cc6d7760aab6c4f95ae1a650ce3b","private":false,"scripts":{"lint":"npx @biomejs/biome check --write","test":"echo \"Error: no test specified\" && exit 1","build":"tsc && tsc-alias"},"_npmUser":{"name":"asammarco","email":"asammarco@hotmail.com"},"_npmVersion":"10.9.0","description":"A reusable MongoDB Change Stream manager","directories":{},"_nodeVersion":"23.2.0","dependencies":{"mongodb":"^6.19.0","typescript":"^5.9.2"},"_hasShrinkwrap":false,"devDependencies":{"tsc-alias":"^1.8.16","@types/node":"^24.3.0","@biomejs/biome":"2.2.2"},"peerDependencies":{"mongodb":"^6.19.0"},"_npmOperationalInternal":{"tmp":"tmp/change_stream_broker_1.0.6_1756508836292_0.6086636005141226","host":"s3://npm-registry-packages-npm-production"},"deprecated":"Package no longer supported. Contact Support at https://www.npmjs.com/support for more info."},"1.0.7":{"name":"@dafaz/change_stream_broker","version":"1.0.7","keywords":["mongodb","change-stream","resume-token","event-driven","microservices"],"author":{"name":"asammarco"},"license":"MIT","_id":"@dafaz/change_stream_broker@1.0.7","maintainers":[{"name":"asammarco","email":"asammarco@hotmail.com"}],"dist":{"shasum":"3e57c3d71e33ead1de578efc4255cdd26e80c965","tarball":"https://registry.npmjs.org/@dafaz/change_stream_broker/-/change_stream_broker-1.0.7.tgz","fileCount":50,"integrity":"sha512-T0aF3Ha6ZQOvGPlaw19AhnwP6f5jurALYqHvPJRoJNKKIlVS0o9SUf3QHw4Ikv1LMlXhFXgx7Ay105N9SpD0fQ==","signatures":[{"sig":"MEYCIQCqEjmY/ExpYc25zt9Hgd2vzbwTdCC5qoQe5YL4TDmLYwIhAMngPlFYNO/KJdTfP2QVmXhTd6VnYaManTEdldJUh2AB","keyid":"SHA256:DhQ8wR5APBvFHLF/+Tc+AYvPOdTpcIDqOhxsBHRwC7U"}],"unpackedSize":124289},"main":"dist/index.js","types":"dist/index.d.ts","gitHead":"f36b9c0613def12827bbd3ee3f74efb292b8e286","private":false,"scripts":{"lint":"npx @biomejs/biome check --write","test":"echo \"Error: no test specified\" && exit 1","build":"tsc && tsc-alias"},"_npmUser":{"name":"asammarco","email":"asammarco@hotmail.com"},"_npmVersion":"10.9.0","description":"A reusable MongoDB Change Stream manager","directories":{},"_nodeVersion":"23.2.0","dependencies":{"mongodb":"^6.19.0","typescript":"^5.9.2"},"_hasShrinkwrap":false,"devDependencies":{"tsc-alias":"^1.8.16","@types/node":"^24.3.0","@biomejs/biome":"2.2.2"},"peerDependencies":{"mongodb":"^6.19.0"},"_npmOperationalInternal":{"tmp":"tmp/change_stream_broker_1.0.7_1756592010896_0.23097148515031107","host":"s3://npm-registry-packages-npm-production"},"deprecated":"Package no longer supported. Contact Support at https://www.npmjs.com/support for more info."},"1.0.8":{"name":"@dafaz/change_stream_broker","version":"1.0.8","keywords":["mongodb","change-stream","resume-token","event-driven","microservices"],"author":{"name":"asammarco"},"license":"MIT","_id":"@dafaz/change_stream_broker@1.0.8","maintainers":[{"name":"asammarco","email":"asammarco@hotmail.com"}],"dist":{"shasum":"94c968f9e83e0c782ea19a0a26fd3e43656f3881","tarball":"https://registry.npmjs.org/@dafaz/change_stream_broker/-/change_stream_broker-1.0.8.tgz","fileCount":62,"integrity":"sha512-2kGIN2vYpl1vdXn+OlkYMQaldUmeYLq9Ma7UA41cjRKw3n+Syp9kCU1OygLC6w6F7aXfLwhAIitIJO0JMfIpLQ==","signatures":[{"sig":"MEMCIBclv0E8d97vKyjoajSSgi4R8KeOCQfF+0PRlJ++JH1GAh8B3Kn4cvchkF8ju1Q6Ohy0w8cjnsTSzLlFl783SZrK","keyid":"SHA256:DhQ8wR5APBvFHLF/+Tc+AYvPOdTpcIDqOhxsBHRwC7U"}],"unpackedSize":162550},"main":"dist/index.js","types":"dist/index.d.ts","gitHead":"5d183f658b213a1b2809cae968c3722e40606125","private":false,"scripts":{"lint":"npx @biomejs/biome check --write","test":"echo \"Error: no test specified\" && exit 1","build":"tsc && tsc-alias"},"_npmUser":{"name":"asammarco","email":"asammarco@hotmail.com"},"_npmVersion":"10.9.0","description":"A reusable MongoDB Change Stream manager","directories":{},"_nodeVersion":"23.2.0","dependencies":{"mongodb":"^6.19.0","typescript":"^5.9.2"},"_hasShrinkwrap":false,"devDependencies":{"tsc-alias":"^1.8.16","@types/node":"^24.3.0","@biomejs/biome":"2.2.2"},"peerDependencies":{"mongodb":"^6.19.0"},"_npmOperationalInternal":{"tmp":"tmp/change_stream_broker_1.0.8_1756593821194_0.11897339393152606","host":"s3://npm-registry-packages-npm-production"},"deprecated":"Package no longer supported. Contact Support at https://www.npmjs.com/support for more info."},"1.0.9":{"name":"@dafaz/change_stream_broker","version":"1.0.9","keywords":["mongodb","change-stream","resume-token","event-driven","microservices"],"author":{"name":"asammarco"},"license":"MIT","_id":"@dafaz/change_stream_broker@1.0.9","maintainers":[{"name":"asammarco","email":"asammarco@hotmail.com"}],"dist":{"shasum":"b6765affc2dca76a2a98044bda58ecbb02eadb56","tarball":"https://registry.npmjs.org/@dafaz/change_stream_broker/-/change_stream_broker-1.0.9.tgz","fileCount":62,"integrity":"sha512-uGY1tiYj0pD6cKXNbwUuAPL6d26HO2lRBtQA1s26jQi476HnhzkMIxJ4mu7m3FRxb46KWexRxp1DtA8GB1+4jA==","signatures":[{"sig":"MEUCID9IZrIYklZVM9OILAYKbPUOoJN65phAZZv1tyRgfyzcAiEAieaPHIBfNdGu7BVhdBVJkEESWFcq6m3bRfvpF4j2Hdg=","keyid":"SHA256:DhQ8wR5APBvFHLF/+Tc+AYvPOdTpcIDqOhxsBHRwC7U"}],"unpackedSize":162692},"main":"dist/index.js","types":"dist/index.d.ts","gitHead":"d6bfc9113b65d8ab74dbe95ad5ede637fb2491aa","private":false,"scripts":{"lint":"npx @biomejs/biome check --write","test":"echo \"Error: no test specified\" && exit 1","build":"tsc && tsc-alias"},"_npmUser":{"name":"asammarco","email":"asammarco@hotmail.com"},"_npmVersion":"10.9.0","description":"A reusable MongoDB Change Stream manager","directories":{},"_nodeVersion":"23.2.0","dependencies":{"mongodb":"^6.19.0","typescript":"^5.9.2"},"_hasShrinkwrap":false,"devDependencies":{"tsc-alias":"^1.8.16","@types/node":"^24.3.0","@biomejs/biome":"2.2.2"},"peerDependencies":{"mongodb":"^6.19.0"},"_npmOperationalInternal":{"tmp":"tmp/change_stream_broker_1.0.9_1756739139446_0.6071843445654901","host":"s3://npm-registry-packages-npm-production"},"deprecated":"Package no longer supported. Contact Support at https://www.npmjs.com/support for more info."},"1.1.0":{"name":"@dafaz/change_stream_broker","version":"1.1.0","keywords":["mongodb","change-stream","resume-token","event-driven","microservices"],"author":{"name":"asammarco"},"license":"MIT","_id":"@dafaz/change_stream_broker@1.1.0","maintainers":[{"name":"asammarco","email":"asammarco@hotmail.com"}],"dist":{"shasum":"3ea41cb3af7204e3f4326e489dc9a3617c4216ed","tarball":"https://registry.npmjs.org/@dafaz/change_stream_broker/-/change_stream_broker-1.1.0.tgz","fileCount":62,"integrity":"sha512-hsLw6rkQ53rzZGOorU7VKp2Ga2X5sfJV59IuwwT5xNxwOrzFUKzylppbcnRm0tsWUNJLm/y+lvEivKPpXkthkQ==","signatures":[{"sig":"MEUCIQCfRHfGSfWia0tsA3V78S6o14lOD2YzHtNJ82pzKDpHsQIgJo6qCCuX2dFs+/zZmdEpa37oS/h5Xg+Rfnyof3DvIHI=","keyid":"SHA256:DhQ8wR5APBvFHLF/+Tc+AYvPOdTpcIDqOhxsBHRwC7U"}],"unpackedSize":203935},"main":"dist/index.js","types":"dist/index.d.ts","gitHead":"07408a7a968c57ac05eeb65d6b377e0fb3a27f89","private":false,"scripts":{"lint":"npx @biomejs/biome check --write","test":"echo \"Error: no test specified\" && exit 1","build":"tsc && tsc-alias"},"_npmUser":{"name":"asammarco","email":"asammarco@hotmail.com"},"_npmVersion":"10.9.0","description":"A reusable MongoDB Change Stream manager","directories":{},"_nodeVersion":"23.2.0","dependencies":{"mongodb":"^6.19.0","typescript":"^5.9.2"},"_hasShrinkwrap":false,"devDependencies":{"tsc-alias":"^1.8.16","@types/node":"^24.3.0","@biomejs/biome":"2.2.2"},"peerDependencies":{"mongodb":"^6.19.0"},"_npmOperationalInternal":{"tmp":"tmp/change_stream_broker_1.1.0_1756852851514_0.12543418631741865","host":"s3://npm-registry-packages-npm-production"},"deprecated":"Package no longer supported. Contact Support at https://www.npmjs.com/support for more info."}},"time":{"created":"2025-08-29T18:47:19.560Z","modified":"2025-09-06T22:22:37.547Z","1.0.0":"2025-08-29T18:47:19.812Z","1.0.1":"2025-08-29T20:45:24.166Z","1.0.2":"2025-08-29T20:55:05.349Z","1.0.3":"2025-08-29T21:02:04.836Z","1.0.4":"2025-08-29T21:16:51.523Z","1.0.5":"2025-08-29T23:05:15.454Z","1.0.6":"2025-08-29T23:07:16.479Z","1.0.7":"2025-08-30T22:13:31.087Z","1.0.8":"2025-08-30T22:43:41.374Z","1.0.9":"2025-09-01T15:05:39.619Z","1.1.0":"2025-09-02T22:40:51.725Z"},"author":{"name":"asammarco"},"license":"MIT","keywords":["mongodb","change-stream","resume-token","event-driven","microservices"],"description":"A reusable MongoDB Change Stream manager","maintainers":[{"name":"asammarco","email":"asammarco@hotmail.com"}],"readme":"# Change Stream Broker\r\n\r\n## 🌐 Português\r\n\r\n### O que é o Projeto\r\n\r\n**Change Stream Broker** é um pacote Node.js que transforma MongoDB Change Streams em um sistema completo de message broker para arquiteturas de microsserviços. Ele fornece uma API similar ao Kafka/RabbitMQ, mas utilizando apenas MongoDB como backbone.\r\n\r\n#### Principais Características\r\n\r\n- 🎯 **API Familiar**: Interface similar a brokers populares (Kafka-like).\r\n- 📦 **Zero Dependências Extras**: Usa apenas o driver oficial do MongoDB.\r\n- 🔄 **Resume Tokens**: Garante delivery exactly-once com mecanismo de offsets.\r\n- 👥 **Consumer Groups**: Suporte a grupos de consumidores com balanceamento automático.\r\n- 🛡️ **Type Safety**: Completo suporte a TypeScript com generics.\r\n- ⚡ **Alta Disponibilidade**: Reconexão automática e lógica de retries integrada.\r\n- 🔧 **Extensível**: Arquitetura modular para customizações.\r\n\r\n#### Casos de Uso Ideais\r\n\r\n- Microsserviços que já utilizam MongoDB.\r\n- Sistemas que precisam de comunicação assíncrona entre serviços.\r\n- Migração de sistemas legados para arquitetura orientada a eventos.\r\n- Ambientes onde Kafka ou RabbitMQ seriam excessivos (overkill).\r\n\r\n### Como Instalar\r\n\r\n#### Pré-requisitos\r\n\r\n- Node.js 16+  \r\n- MongoDB 5.0+ com replica set habilitado  **(veja exemplo no final)**\r\n- TypeScript **(recomendado)**\r\n\r\n#### Instalação\r\n\r\n```bash\r\nnpm install @dafaz/change-stream-broker\r\n```\r\n\r\n#### Configuração do MongoDB\r\n\r\nPara usar **Change Streams**, é necessário configurar um replica set:\r\n\r\n- Utilize uma instância do **MongoDB Atlas Online**, que já vem com replica set habilitado.\r\n- Ou configure um replica set localmente utilizando um container Docker. Consulte a documentação oficial do MongoDB para mais detalhes.\r\n\r\n### Estrutura de Arquivos do Projeto\r\n\r\n#### Principais Classes e Interfaces\r\n\r\n```text\r\nsrc/\r\n├── index.ts                          # Ponto de entrada principal\r\n├── broker/\r\n│   ├── change-stream-broker.ts       # Classe principal do broker\r\n│   ├── topic-manager.ts              # Gerenciador de tópicos\r\n│   └── types.ts                      # Interfaces principais\r\n├── consumer/\r\n│   ├── consumer.ts                   # Implementação do consumer\r\n│   ├── consumer.interface.ts         # Interface do consumer\r\n│   ├── consumer-group.ts             # Gerenciador de consumer groups\r\n│   └── message-handler.interface.ts  # Interface de handlers\r\n├── producer/\r\n│   └── producer.ts                   # Implementação do producer\r\n└── storage/\r\n    ├── offset-storage.ts             # Interface de storage\r\n    ├── mongo-offset-storage.ts       # Storage de offsets no MongoDB\r\n    └── file-offset-storage.ts        # Storage de offsets em arquivo\r\n```\r\n\r\n### Principais Interfaces\r\n\r\n#### IChangeStreamConsumer\r\n\r\n- **Definição**:  \r\n```typescript\r\n  interface IChangeStreamConsumer {\r\n  getConsumerId(): string;\r\n  getGroupId(): string;\r\n  getTopic(): string;\r\n  connect(): Promise<void>;\r\n  subscribe<T = any>(config: MessageHandlerConfig<T>): Promise<void>;\r\n  disconnect(): Promise<void>;\r\n  commitOffsets(): Promise<void>;\r\n}\r\n```\r\nA interface `IChangeStreamConsumer` define o contrato para a implementação de um consumidor de mensagens no **Change Stream Broker**. Ela garante que qualquer classe que implemente essa interface possua os métodos necessários para gerenciar a conexão, assinatura e processamento de mensagens. Abaixo estão os métodos definidos:\r\n\r\n- **`getConsumerId(): string`**  \r\n  Retorna o identificador único do consumidor. Esse ID é usado para distinguir consumidores individuais dentro de um grupo.\r\n\r\n- **`getGroupId(): string`**  \r\n  Retorna o identificador do grupo de consumidores ao qual este consumidor pertence. Consumidores no mesmo grupo compartilham a carga de processamento de mensagens.\r\n\r\n- **`getTopic(): string`**  \r\n  Retorna o nome do tópico ao qual o consumidor está associado. O tópico é a fonte das mensagens consumidas.\r\n\r\n- **`connect(): Promise<void>`**  \r\n  Estabelece a conexão do consumidor com o broker. Este método deve ser chamado antes de qualquer operação de assinatura ou consumo.\r\n\r\n- **`subscribe<T = any>(config: MessageHandlerConfig<T>): Promise<void>`**  \r\n  Permite que o consumidor se inscreva em um tópico com uma configuração específica de manipulador de mensagens (`MessageHandlerConfig`). O tipo genérico `<T>` define o formato esperado das mensagens.\r\n\r\n- **`disconnect(): Promise<void>`**  \r\n  Desconecta o consumidor do broker, encerrando a assinatura e liberando os recursos associados.\r\n\r\n- **`commitOffsets(): Promise<void>`**  \r\n  Confirma os offsets das mensagens processadas, garantindo que elas não sejam reprocessadas em caso de falhas ou reinicializações.\r\n\r\nEssa interface é essencial para garantir que os consumidores sigam um padrão consistente, facilitando a implementação e manutenção do sistema.\r\n\r\n\r\n#### IChangeStreamProducer\r\n\r\n```typescript\r\nexport interface IChangeStreamProducer {\r\n  send(messages: Message | Message[]): Promise<void>\r\n  disconnect(): Promise<void>\r\n  isConnected: boolean\r\n}\r\n```\r\nMétodos:\r\n\r\n- **`send(messages: Message | Message[]): Promise<void>`**  \r\nEnvia uma ou mais mensagens para o tópico\r\n\r\n- **`disconnect(): Promise<void>`**  \r\nDesconecta o producer\r\n\r\n- **`isConnected(): boolean`**  \r\nVerifica se está conectado\r\n\r\n\r\n#### MessageHandler\r\n\r\nA interface `MessageHandler` define o contrato para a função responsável por processar mensagens consumidas de um tópico. Essa função é utilizada pelos consumidores para lidar com cada mensagem recebida. Abaixo está a explicação detalhada:\r\n\r\n- **Definição**:  \r\n  ```typescript\r\n  interface MessageHandler<T = unknown> {\r\n    (record: ConsumerRecord & { message: { value: T } }): Promise<void>;\r\n  }\r\n  ```\r\nEssa interface é essencial para definir como as mensagens devem ser processadas pelos consumidores, permitindo que cada mensagem seja tratada de forma personalizada e assíncrona.\r\n\r\nTipo Genérico:\r\n- **`<T = unknown>`**  \r\nPermite que o manipulador seja configurado para processar mensagens de diferentes formatos, garantindo flexibilidade e segurança de tipo.\r\n\r\nParâmetros:\r\n- **`record`**  \r\n  Um objeto que combina as propriedades de ConsumerRecord com uma mensagem (message) contendo um valor (value) do tipo genérico T.\r\n\r\n- **`ConsumerRecord`**  \r\n  Representa os metadados da mensagem consumida, como o tópico, partição e offset.\r\n\r\n - **`message.value`**  \r\n  Contém o valor da mensagem, que é do tipo genérico T.\r\n\r\n\r\n#### ConsumerRecord\r\n\r\nA interface `ConsumerRecord` representa os metadados de uma mensagem consumida de um tópico. Esses metadados fornecem informações importantes sobre a origem e o contexto da mensagem, além de dados necessários para gerenciar o consumo. Abaixo está a explicação detalhada:\r\n\r\n- **Definição**:  \r\n  ```typescript\r\n  interface ConsumerRecord {\r\n    topic: string;\r\n    partition: number;\r\n    message: Message;\r\n    offset: ResumeToken;\r\n    timestamp: Date;\r\n  }\r\n  ```\r\nEssa interface é fundamental para fornecer o contexto completo de uma mensagem consumida, permitindo que os consumidores processem as mensagens de forma eficiente e segura.\r\n\r\nPropriedades:\r\n- **`topic: string`**  \r\nO nome do tópico de onde a mensagem foi consumida.\r\n\r\n- **`partition: number`**  \r\nO número da partição do tópico de onde a mensagem foi lida. Partições são usadas para distribuir mensagens e balancear a carga entre consumidores.\r\n\r\n- **`message: Message`**  \r\nO conteúdo da mensagem consumida. A interface Message contém os dados principais que serão processados.\r\n\r\n- **`offset: ResumeToken`**  \r\nO token de resume do MongoDB, usado para rastrear a posição da mensagem no stream. Esse token é essencial para garantir o processamento exatamente uma vez (exactly-once) e para retomar o consumo em caso de falhas.\r\n\r\n- **`timestamp: Date`**  \r\nA data e hora em que a mensagem foi produzida ou consumida. Esse valor pode ser usado para fins de auditoria ou ordenação.\r\n\r\n\r\n#### OffsetStorage\r\n\r\nA interface `OffsetStorage` define o contrato para o gerenciamento de offsets no **Change Stream Broker**. Ela é responsável por armazenar e recuperar os offsets das mensagens consumidas, garantindo que o sistema possa retomar o consumo de onde parou em caso de falhas ou reinicializações. Abaixo está a explicação detalhada:\r\n\r\n- **Definição**:  \r\n  ```typescript\r\n  interface OffsetStorage {\r\n    commitOffset(commit: OffsetCommit): Promise<void>;\r\n    getOffset(groupId: string, topic: string, partition: number): Promise<ResumeToken | null>;\r\n  }\r\n  ```\r\nEssa interface é essencial para implementar diferentes estratégias de armazenamento de offsets, como armazenamento em MongoDB, arquivos ou outros sistemas, garantindo flexibilidade e confiabilidade no gerenciamento do consumo de mensagens.\r\n\r\nMétodo:\r\n- **`commitOffset(commit: OffsetCommit): Promise<void>`**  \r\nArmazena o offset mais recente para um grupo de consumidores, tópico e partição específicos.\r\n\r\nParâmetro:\r\n- **`commit: OffsetCommit`**  \r\nUm objeto contendo as informações necessárias para registrar o offset, como o grupo de consumidores, tópico, partição e o token de resume.\r\n\r\nMétodo:\r\n- **`getOffset(groupId: string, topic: string, partition: number): Promise<ResumeToken | null>`**  \r\nRecupera o offset armazenado para um grupo de consumidores, tópico e partição específicos.\r\n\r\nParâmetros:\r\n- **`groupId: string`**  \r\nO identificador do grupo de consumidores.\r\n\r\n- **`topic: string`**  \r\nO nome do tópico.\r\n\r\n- **`partition: number`**  \r\nO número da partição.\r\n\r\nRetorno:\r\n- **`Promise<ResumeToken | null>`**  \r\nRetorno assíncrono do token de resume armazenado ou null se nenhum offset estiver disponível.\r\n\r\n\r\n### Fluxo de Dados\r\n\r\n```text\r\n[Producer Service]          [MongoDB]               [Consumer Service]\r\n     |                         |                         |\r\n     | 1. Insere documento     |                         |\r\n     |-----------------------> |                         |\r\n     |                         |                         |\r\n     |                         | 2. Change Stream detecta|\r\n     |                         |    mudança              |\r\n     |                         |-----------------------> |\r\n     |                         |                         |\r\n     |                         | 3. Processa mensagem    |\r\n     |                         |    e commit offset      |\r\n     |                         |<----------------------- |\r\n```\r\n\r\n### Configurações de Performance\r\n\r\nO exemplo abaixo demonstra como configurar o **Change Stream Broker** com parâmetros otimizados para garantir alta performance e resiliência:\r\n\r\n```typescript\r\n// Exemplo de configuração otimizada\r\nconst broker = new ChangeStreamBroker({\r\n  mongoUri: 'mongodb://localhost:27017', // URI de conexão com o MongoDB\r\n  database: 'my-events',                // Nome do banco de dados utilizado\r\n  maxRetries: 10,                       // Número máximo de tentativas de reconexão em caso de falhas\r\n  retryDelayMs: 1000,                   // Intervalo (em milissegundos) entre as tentativas de reconexão\r\n  heartbeatIntervalMs: 30000,           // Intervalo (em milissegundos) para envio de heartbeats\r\n  autoCreateTopics: false               // Criação automática de tópicos desabilitada para evitar erro\r\n});\r\n```\r\nEssa configuração é ideal para cenários onde a resiliência e a estabilidade são cruciais, garantindo que o sistema continue funcionando mesmo em situações de falhas temporárias.\r\n\r\nParâmetros:  \r\n- **`mongoUri`**  \r\nDefine a URI de conexão com o MongoDB. No exemplo, está configurado para um MongoDB local na porta padrão (27017).\r\n\r\n- **`database`**  \r\nEspecifica o banco de dados onde os eventos serão armazenados e gerenciados.\r\n\r\n- **`maxRetries`**  \r\nConfigura o número máximo de tentativas de reconexão em caso de falhas na comunicação com o MongoDB.\r\n\r\n- **`retryDelayMs`**  \r\nDefine o tempo de espera (em milissegundos) entre cada tentativa de reconexão.\r\n\r\n- **`heartbeatIntervalMs`**  \r\nDetermina o intervalo de envio de heartbeats para monitorar a conexão com o MongoDB.\r\n\r\n- **`autoCreateTopics`**  \r\nQuando habilitado (true), permite que o broker crie tópicos automaticamente, caso eles não existam.\r\n\r\n\r\n### Exemplo de Microsserviços com NestJS e Change Stream Broker\r\n\r\n#### Estrutura do Projeto\r\n\r\n```text\r\nmicroservices-example/\r\n├── purchase-service/\r\n│   ├── src/\r\n│   │   ├── purchases/\r\n│   │   │   ├── purchases.module.ts\r\n│   │   │   ├── purchases.service.ts\r\n│   │   │   ├── purchases.controller.ts\r\n│   │   │   └── schemas/\r\n│   │   │       └── purchase.schema.ts\r\n│   │   ├── change-stream/\r\n│   │   │   └── purchase.publisher.ts\r\n│   │   ├── app.module.ts\r\n│   │   └── main.ts\r\n│   ├── package.json\r\n│   └── Dockerfile\r\n├── classroom-service/\r\n│   ├── src/\r\n│   │   ├── enrollments/\r\n│   │   │   ├── enrollments.module.ts\r\n│   │   │   ├── enrollments.service.ts\r\n│   │   │   ├── enrollments.controller.ts\r\n│   │   │   └── schemas/\r\n│   │   │       └── enrollment.schema.ts\r\n│   │   ├── change-stream/\r\n│   │   │   └── enrollment.consumer.ts\r\n│   │   ├── app.module.ts\r\n│   │   └── main.ts\r\n│   ├── package.json\r\n│   └── Dockerfile\r\n├── docker-compose.yml\r\n└── package.json\r\n```\r\n\r\n### Serviço de Purchase (Producer)\r\n\r\n\r\n#### package.json\r\n\r\n**purchase-service/ package.json**  \r\n\r\n```json\r\n{\r\n  \"name\": \"purchase-service\",\r\n  \"version\": \"1.0.0\",\r\n  \"dependencies\": {\r\n    \"@nestjs/common\": \"^9.0.0\",\r\n    \"@nestjs/core\": \"^9.0.0\",\r\n    \"@nestjs/mongoose\": \"^9.0.0\",\r\n    \"@nestjs/platform-express\": \"^9.0.0\",\r\n    \"mongoose\": \"^6.0.0\",\r\n    \"@dafaz/change-stream-broker\": \"^1.0.0\"\r\n  }\r\n}\r\n```\r\n\r\n#### PurchaseSchema\r\n\r\n**purchase-service/ src/ purchases/ schemas/ purchase.schema.ts**  \r\n\r\n```typescript\r\nimport { Prop, Schema, SchemaFactory } from '@nestjs/mongoose';\r\nimport { Document, Schema as MongooseSchema } from 'mongoose';\r\n\r\nexport type PurchaseDocument = Purchase & Document;\r\n\r\n@Schema({ timestamps: true })\r\nexport class Purchase {\r\n  @Prop({ required: true, type: MongooseSchema.Types.ObjectId, ref: 'Customer' })\r\n  customerId: string;\r\n\r\n  @Prop({ required: true, type: MongooseSchema.Types.ObjectId, ref: 'Product' })\r\n  productId: string;\r\n\r\n  @Prop({ required: true })\r\n  productType: string;\r\n\r\n  @Prop({ required: true })\r\n  amount: number;\r\n\r\n  @Prop({ default: 'pending' })\r\n  status: string;\r\n\r\n  @Prop()\r\n  transactionId: string;\r\n}\r\n\r\nexport const PurchaseSchema = SchemaFactory.createForClass(Purchase);\r\n```\r\n\r\n#### PurchasePublisher\r\n\r\n**purchase-service/ src/ change-stream/ purchase.publisher.ts**  \r\n\r\n```typescript\r\nimport { Injectable, OnModuleInit, OnModuleDestroy } from '@nestjs/common';\r\nimport { \r\n  ChangeStreamBroker, \r\n  ProducerConfig,\r\n  IChangeStreamProducer  // ← Importar a interface\r\n} from '@dafaz/change-stream-broker';\r\n\r\nexport interface PurchaseMessage {\r\n  purchaseId: string;\r\n  customerId: string;\r\n  productId: string;\r\n  productType: string;\r\n  amount: number;\r\n  status: string;\r\n  createdAt: Date;\r\n}\r\n\r\n@Injectable()\r\nexport class PurchasePublisher implements OnModuleInit, OnModuleDestroy {\r\n  private broker: ChangeStreamBroker;\r\n  private producer: IChangeStreamProducer;  // ← Usando interface\r\n\r\n  constructor() {\r\n    this.broker = new ChangeStreamBroker({\r\n      mongoUri: process.env.MONGO_URI || 'mongodb://localhost:27017',\r\n      database: 'purchase-events',\r\n      autoCreateTopics: true\r\n    });\r\n  }\r\n\r\n  async onModuleInit() {\r\n    await this.broker.connect();\r\n    \r\n    const producerConfig: ProducerConfig = {\r\n      topic: 'purchases'\r\n    };\r\n    \r\n    this.producer = await this.broker.createProducer(producerConfig);\r\n    console.log('Purchase publisher initialized');\r\n  }\r\n\r\n  async publishPurchaseCreated(purchase: PurchaseMessage) {\r\n    await this.producer.send({\r\n      key: purchase.purchaseId,\r\n      value: purchase,\r\n      headers: {\r\n        'event-type': 'purchase.created',\r\n        'source': 'purchase-service'\r\n      }\r\n    });\r\n    \r\n    console.log(`Purchase event published: ${purchase.purchaseId}`);\r\n  }\r\n\r\n  async publishPurchaseStatusChanged(purchaseId: string, status: string) {\r\n    await this.producer.send({\r\n      key: purchaseId,\r\n      value: { purchaseId, status, updatedAt: new Date() },\r\n      headers: {\r\n        'event-type': 'purchase.status.changed',\r\n        'source': 'purchase-service'\r\n      }\r\n    });\r\n  }\r\n\r\n  async onModuleDestroy() {\r\n    await this.broker.disconnect();\r\n  }\r\n}\r\n```\r\n\r\n#### PurchasesService\r\n\r\n**purchase-service/ src/ purchases/ purchases.service.ts**  \r\n\r\n```typescript\r\nimport { Injectable } from '@nestjs/common';\r\nimport { InjectModel } from '@nestjs/mongoose';\r\nimport { Model } from 'mongoose';\r\nimport { Purchase, PurchaseDocument } from './schemas/purchase.schema';\r\nimport { PurchasePublisher } from '../change-stream/purchase.publisher';\r\n\r\n@Injectable()\r\nexport class PurchasesService {\r\n  constructor(\r\n    @InjectModel(Purchase.name) private purchaseModel: Model<PurchaseDocument>,\r\n    private purchasePublisher: PurchasePublisher\r\n  ) {}\r\n\r\n  async createPurchase(createPurchaseDto: any) {\r\n    const purchase = new this.purchaseModel(createPurchaseDto);\r\n    const savedPurchase = await purchase.save();\r\n\r\n    // Publicar evento de purchase created\r\n    await this.purchasePublisher.publishPurchaseCreated({\r\n      purchaseId: savedPurchase._id.toString(),\r\n      customerId: savedPurchase.customerId,\r\n      productId: savedPurchase.productId,\r\n      productType: savedPurchase.productType,\r\n      amount: savedPurchase.amount,\r\n      status: savedPurchase.status,\r\n      createdAt: savedPurchase.createdAt\r\n    });\r\n\r\n    return savedPurchase;\r\n  }\r\n\r\n  async updatePurchaseStatus(purchaseId: string, status: string) {\r\n    const updatedPurchase = await this.purchaseModel.findByIdAndUpdate(\r\n      purchaseId,\r\n      { status },\r\n      { new: true }\r\n    );\r\n\r\n    if (updatedPurchase) {\r\n      await this.purchasePublisher.publishPurchaseStatusChanged(\r\n        purchaseId,\r\n        status\r\n      );\r\n    }\r\n\r\n    return updatedPurchase;\r\n  }\r\n\r\n  async findByCustomerId(customerId: string) {\r\n    return this.purchaseModel.find({ customerId });\r\n  }\r\n}\r\n```\r\n\r\n#### PurchasesController\r\n\r\n**purchase-service/ src/ purchases/ purchases.controller.ts**  \r\n\r\n```typescript\r\nimport { Controller, Post, Body, Param, Patch, Get } from '@nestjs/common';\r\nimport { PurchasesService } from './purchases.service';\r\n\r\n@Controller('purchases')\r\nexport class PurchasesController {\r\n  constructor(private readonly purchasesService: PurchasesService) {}\r\n\r\n  @Post()\r\n  async create(@Body() createPurchaseDto: any) {\r\n    return this.purchasesService.createPurchase(createPurchaseDto);\r\n  }\r\n\r\n  @Patch(':id/status')\r\n  async updateStatus(@Param('id') id: string, @Body('status') status: string) {\r\n    return this.purchasesService.updatePurchaseStatus(id, status);\r\n  }\r\n\r\n  @Get('customer/:customerId')\r\n  async findByCustomer(@Param('customerId') customerId: string) {\r\n    return this.purchasesService.findByCustomerId(customerId);\r\n  }\r\n}\r\n```\r\n\r\n#### AppModule\r\n\r\n**purchase-service/ src/ app.module.ts**  \r\n\r\n```typescript\r\nimport { Module } from '@nestjs/common';\r\nimport { MongooseModule } from '@nestjs/mongoose';\r\nimport { PurchasesModule } from './purchases/purchases.module';\r\nimport { PurchasePublisher } from './change-stream/purchase.publisher';\r\n\r\n@Module({\r\n  imports: [\r\n    MongooseModule.forRoot(process.env.MONGO_URI || 'mongodb://localhost:27017/purchase-service'),\r\n    PurchasesModule,\r\n  ],\r\n  providers: [PurchasePublisher],\r\n  exports: [PurchasePublisher],\r\n})\r\nexport class AppModule {}\r\n```\r\n\r\n### Serviço de Classroom (Consumer)\r\n\r\n\r\n#### package.json\r\n\r\n**classroom-service/ package.json**  \r\n\r\n```json\r\n{\r\n  \"name\": \"classroom-service\",\r\n  \"version\": \"1.0.0\",\r\n  \"dependencies\": {\r\n    \"@nestjs/common\": \"^9.0.0\",\r\n    \"@nestjs/core\": \"^9.0.0\",\r\n    \"@nestjs/mongoose\": \"^9.0.0\",\r\n    \"@nestjs/platform-express\": \"^9.0.0\",\r\n    \"mongoose\": \"^6.0.0\",\r\n    \"@dafaz/change-stream-broker\": \"^1.0.0\"\r\n  }\r\n}\r\n```\r\n\r\n\r\n#### EnrollmentSchema\r\n\r\n**classroom-service/ src/ enrollments/ schemas/ enrollment.schema.ts**  \r\n\r\n```typescript\r\nimport { Prop, Schema, SchemaFactory } from '@nestjs/mongoose';\r\nimport { Document, Schema as MongooseSchema } from 'mongoose';\r\n\r\nexport type EnrollmentDocument = Enrollment & Document;\r\n\r\n@Schema({ timestamps: true })\r\nexport class Enrollment {\r\n  @Prop({ required: true, type: MongooseSchema.Types.ObjectId, ref: 'Student' })\r\n  studentId: string;\r\n\r\n  @Prop({ required: true, type: MongooseSchema.Types.ObjectId, ref: 'Course' })\r\n  courseId: string;\r\n\r\n  @Prop({ required: true, unique: true })\r\n  purchaseId: string;\r\n\r\n  @Prop({ default: 'active' })\r\n  status: string;\r\n\r\n  @Prop()\r\n  enrolledAt: Date;\r\n\r\n  @Prop()\r\n  completedAt: Date;\r\n}\r\n\r\nexport const EnrollmentSchema = SchemaFactory.createForClass(Enrollment);\r\n```\r\n\r\n\r\n#### PurchaseMessage\r\n\r\n**classroom-service/ src/ change-stream/ enrollment.consumer.ts**  \r\n\r\n```typescript\r\nimport { Injectable, OnModuleInit, OnModuleDestroy } from '@nestjs/common';\r\nimport { \r\n  ChangeStreamBroker, \r\n  ConsumerConfig, \r\n  MessageHandlerConfig \r\n} from '@dafaz/change-stream-broker';\r\nimport { EnrollmentsService } from '../enrollments/enrollments.service';\r\n\r\ninterface PurchaseMessage {\r\n  purchaseId: string;\r\n  customerId: string;\r\n  productId: string;\r\n  productType: string;\r\n  amount: number;\r\n  status: string;\r\n  createdAt: Date;\r\n}\r\n\r\n@Injectable()\r\nexport class EnrollmentConsumer implements OnModuleInit, OnModuleDestroy {\r\n  private broker: ChangeStreamBroker;\r\n\r\n  constructor(private enrollmentsService: EnrollmentsService) {\r\n    this.broker = new ChangeStreamBroker({\r\n      mongoUri: process.env.MONGO_URI || 'mongodb://localhost:27017',\r\n      database: 'purchase-events'\r\n    });\r\n  }\r\n\r\n  async onModuleInit() {\r\n    await this.broker.connect();\r\n\r\n    const consumerConfig: ConsumerConfig = {\r\n      groupId: 'classroom-service',\r\n      topic: 'purchases',\r\n      autoCommit: true,\r\n      autoCommitIntervalMs: 5000,\r\n      fromBeginning: false\r\n    };\r\n\r\n    const handlerConfig: MessageHandlerConfig<PurchaseMessage> = {\r\n      handler: this.handlePurchaseEvent.bind(this),\r\n      errorHandler: this.handleError.bind(this),\r\n      maxRetries: 3,\r\n      retryDelay: 1000,\r\n      autoCommit: true\r\n    };\r\n\r\n    const consumer = await this.broker.createConsumer(consumerConfig);\r\n    await consumer.subscribe(handlerConfig);\r\n\r\n    console.log('Enrollment consumer started listening for purchase events...');\r\n  }\r\n\r\n  private async handlePurchaseEvent(record: any) {\r\n    const purchase = record.message.value;\r\n    \r\n    // Apenas processar compras de cursos completadas\r\n    if (purchase.productType === 'course' && purchase.status === 'completed') {\r\n      console.log(`Processing course purchase: ${purchase.purchaseId}`);\r\n      \r\n      try {\r\n        // Verificar se já existe matrícula para esta compra\r\n        const existingEnrollment = await this.enrollmentsService.findByPurchaseId(purchase.purchaseId);\r\n        \r\n        if (!existingEnrollment) {\r\n          // Criar matrícula automaticamente\r\n          await this.enrollmentsService.create({\r\n            studentId: purchase.customerId, // customerId é o studentId\r\n            courseId: purchase.productId,   // productId é o courseId\r\n            purchaseId: purchase.purchaseId,\r\n            enrolledAt: new Date()\r\n          });\r\n          \r\n          console.log(`Enrollment created for purchase: ${purchase.purchaseId}`);\r\n        }\r\n      } catch (error) {\r\n        console.error(`Error creating enrollment for purchase ${purchase.purchaseId}:`, error);\r\n        throw error; // Não commitar offset para reprocessar\r\n      }\r\n    }\r\n  }\r\n\r\n  private async handleError(error: Error, record?: any) {\r\n    console.error('Error processing purchase event:', error);\r\n    \r\n    if (record) {\r\n      console.error('Failed purchase record:', record.message.value);\r\n      // Poderia enviar para uma DLQ aqui\r\n    }\r\n  }\r\n\r\n  async onModuleDestroy() {\r\n    await this.broker.disconnect();\r\n  }\r\n}\r\n```\r\n\r\n#### EnrollmentsService\r\n\r\n**classroom-service/ src/ enrollments/ enrollments.service.ts**  \r\n\r\n```typescript\r\nimport { Injectable } from '@nestjs/common';\r\nimport { InjectModel } from '@nestjs/mongoose';\r\nimport { Model } from 'mongoose';\r\nimport { Enrollment, EnrollmentDocument } from './schemas/enrollment.schema';\r\n\r\n@Injectable()\r\nexport class EnrollmentsService {\r\n  constructor(\r\n    @InjectModel(Enrollment.name) private enrollmentModel: Model<EnrollmentDocument>\r\n  ) {}\r\n\r\n  async create(createEnrollmentDto: any) {\r\n    const enrollment = new this.enrollmentModel(createEnrollmentDto);\r\n    return enrollment.save();\r\n  }\r\n\r\n  async findByPurchaseId(purchaseId: string) {\r\n    return this.enrollmentModel.findOne({ purchaseId });\r\n  }\r\n\r\n  async findByStudentId(studentId: string) {\r\n    return this.enrollmentModel.find({ studentId }).populate('courseId');\r\n  }\r\n\r\n  async updateStatus(enrollmentId: string, status: string) {\r\n    return this.enrollmentModel.findByIdAndUpdate(\r\n      enrollmentId,\r\n      { status },\r\n      { new: true }\r\n    );\r\n  }\r\n\r\n  async completeEnrollment(enrollmentId: string) {\r\n    return this.enrollmentModel.findByIdAndUpdate(\r\n      enrollmentId,\r\n      { \r\n        status: 'completed',\r\n        completedAt: new Date()\r\n      },\r\n      { new: true }\r\n    );\r\n  }\r\n}\r\n```\r\n\r\n#### EnrollmentsController\r\n\r\n**classroom-service/ src/ enrollments/ enrollments.controller.ts**  \r\n\r\n```typescript\r\nimport { Controller, Get, Post, Body, Param, Patch } from '@nestjs/common';\r\nimport { EnrollmentsService } from './enrollments.service';\r\n\r\n@Controller('enrollments')\r\nexport class EnrollmentsController {\r\n  constructor(private readonly enrollmentsService: EnrollmentsService) {}\r\n\r\n  @Post()\r\n  async create(@Body() createEnrollmentDto: any) {\r\n    return this.enrollmentsService.create(createEnrollmentDto);\r\n  }\r\n\r\n  @Get('student/:studentId')\r\n  async findByStudent(@Param('studentId') studentId: string) {\r\n    return this.enrollmentsService.findByStudentId(studentId);\r\n  }\r\n\r\n  @Patch(':id/status')\r\n  async updateStatus(@Param('id') id: string, @Body('status') status: string) {\r\n    return this.enrollmentsService.updateStatus(id, status);\r\n  }\r\n\r\n  @Patch(':id/complete')\r\n  async completeEnrollment(@Param('id') id: string) {\r\n    return this.enrollmentsService.completeEnrollment(id);\r\n  }\r\n}\r\n```\r\n\r\n#### AppModule\r\n\r\n**classroom-service/ src/ app.module.ts**\r\n\r\n```typescript\r\nimport { Module } from '@nestjs/common';\r\nimport { MongooseModule } from '@nestjs/mongoose';\r\nimport { EnrollmentsModule } from './enrollments/enrollments.module';\r\nimport { EnrollmentConsumer } from './change-stream/enrollment.consumer';\r\n\r\n@Module({\r\n  imports: [\r\n    MongooseModule.forRoot(process.env.MONGO_URI || 'mongodb://localhost:27017/classroom-service'),\r\n    EnrollmentsModule,\r\n  ],\r\n  providers: [EnrollmentConsumer],\r\n})\r\nexport class AppModule {}\r\n```\r\n\r\n### Exemplo de Configuração Docker\r\n\r\n#### replSet local (para desenvolvimento)\r\n\r\nO Change Stream exige que a sua instância do MongoDB esteja em um cluster.\r\nPara fins de desenvolvimento, você pode simular esse cluster. Vamos criar dois\r\narquivos para subir uma instância de dados no docker com replSet configurado.\r\n**Observação: isso só deve ser feito para fins de desenvolvimento.**\r\n\r\n[root] / docker-compose.yml\r\n\r\n```yml\r\nversion: '3.8'\r\n \r\nservices:\r\n  mongo:\r\n    build:\r\n      context: .\r\n      dockerfile: Dockerfile\r\n    container_name: mongo-server\r\n    restart: always\r\n    environment:\r\n      MONGO_INITDB_ROOT_USERNAME: root\r\n      MONGO_INITDB_ROOT_PASSWORD: docker\r\n    ports:\r\n      - \"0.0.0.0:27017:27017\"\r\n    expose:\r\n      - 27017\r\n    command: --replSet rs0 --keyFile /etc/mongo-keyfile --bind_ip_all --port 27017 --auth\r\n    healthcheck:\r\n      test: echo \"try { rs.status() } catch (err) { rs.initiate({_id:'rs0',members:[{_id:0,host:'127.0.0.1:27017'}]}) }\" | mongosh --port 27017 -u root -p docker --authenticationDatabase admin\r\n      interval: 5s\r\n      timeout: 15s\r\n      start_period: 15s\r\n      retries: 10\r\n    volumes:\r\n      - \"./data:/data/db\"\r\n      - \"./data/config:/etc/mongod.conf\"\r\n    networks:\r\n      - mongo-net\r\n  \r\nnetworks:\r\n  mongo-net:\r\n    driver: bridge\r\n```\r\n\r\n[root-path] / Dockerfile\r\n\r\n```text\r\nFROM mongo\r\nRUN openssl rand -base64 756 > /etc/mongo-keyfile \r\nRUN chmod 400 /etc/mongo-keyfile \r\nRUN chown mongodb:mongodb /etc/mongo-keyfile \r\n```\r\n","readmeFilename":"README.md"}