diff --git a/.github/workflows/build-samples.yml b/.github/workflows/build-samples.yml index 2421bcc5..1a349b30 100644 --- a/.github/workflows/build-samples.yml +++ b/.github/workflows/build-samples.yml @@ -231,6 +231,10 @@ jobs: - name: Install DF - work-item-filtering-app-b run: pip install -r samples/durable-functions/python/work-item-filtering-app-b/requirements.txt + # Scenario samples + - name: Install Scenario - WorkItemFilteringSplitActivitiesPython + run: pip install -r samples/scenarios/WorkItemFilteringSplitActivitiesPython/src/client/requirements.txt + # Syntax check all Python files - name: Syntax check Python files run: find samples -name "*.py" -not -path "*/__pycache__/*" -exec python -m py_compile {} + diff --git a/samples/durable-functions/javascript/HelloCities/package-lock.json b/samples/durable-functions/javascript/HelloCities/package-lock.json new file mode 100644 index 00000000..529f6e15 --- /dev/null +++ b/samples/durable-functions/javascript/HelloCities/package-lock.json @@ -0,0 +1,482 @@ +{ + "name": "durable-functions-hello-cities-js", + "version": "1.0.0", + "lockfileVersion": 3, + "requires": true, + "packages": { + "": { + "name": "durable-functions-hello-cities-js", + "version": "1.0.0", + "dependencies": { + "@azure/functions": "^4.0.0", + "durable-functions": "^3.0.0" + } + }, + "node_modules/@azure/functions": { + "version": "4.16.0", + "resolved": "https://registry.npmjs.org/@azure/functions/-/functions-4.16.0.tgz", + "integrity": "sha512-cwLld7tRi1VLYNdm1evgS9n8ABP/e+0uLjRO++OjM/kigERG80OfFRRMbag/vpXUzCtda83lDiuCvcwPlxgwZw==", + "license": "MIT", + "dependencies": { + "@azure/functions-extensions-base": "0.3.0", + "cookie": "^0.7.0" + }, + "engines": { + "node": ">=20.0" + } + }, + "node_modules/@azure/functions-extensions-base": { + "version": "0.3.0", + "resolved": "https://registry.npmjs.org/@azure/functions-extensions-base/-/functions-extensions-base-0.3.0.tgz", + "integrity": "sha512-Cux0hLu5ZXlC/Kb+yvJVhRLIdkfFwui2HeT5oGZL00r/GCUUkhGTzRfZUjRN4Bq729mPv3okPucz2z7SMQLStA==", + "license": "MIT", + "engines": { + "node": ">=18.0" + } + }, + "node_modules/@opentelemetry/api": { + "version": "1.9.1", + "resolved": "https://registry.npmjs.org/@opentelemetry/api/-/api-1.9.1.tgz", + "integrity": "sha512-gLyJlPHPZYdAk1JENA9LeHejZe1Ti77/pTeFm/nMXmQH/HFZlcS/O2XJB+L8fkbrNSqhdtlvjBVjxwUYanNH5Q==", + "license": "Apache-2.0", + "engines": { + "node": ">=8.0.0" + } + }, + "node_modules/agent-base": { + "version": "6.0.2", + "resolved": "https://registry.npmjs.org/agent-base/-/agent-base-6.0.2.tgz", + "integrity": "sha512-RZNwNclF7+MS/8bDg70amg32dyeZGZxiDuQmZxKLAlQjr3jGyLx+4Kkk58UO7D2QdgFIQCovuSuZESne6RG6XQ==", + "license": "MIT", + "dependencies": { + "debug": "4" + }, + "engines": { + "node": ">= 6.0.0" + } + }, + "node_modules/agent-base/node_modules/debug": { + "version": "4.4.3", + "resolved": "https://registry.npmjs.org/debug/-/debug-4.4.3.tgz", + "integrity": "sha512-RGwwWnwQvkVfavKVt22FGLw+xYSdzARwm0ru6DhTVA3umU5hZc28V3kO4stgYryrTlLpuvgI9GiijltAjNbcqA==", + "license": "MIT", + "dependencies": { + "ms": "^2.1.3" + }, + "engines": { + "node": ">=6.0" + }, + "peerDependenciesMeta": { + "supports-color": { + "optional": true + } + } + }, + "node_modules/agent-base/node_modules/ms": { + "version": "2.1.3", + "resolved": "https://registry.npmjs.org/ms/-/ms-2.1.3.tgz", + "integrity": "sha512-6FlzubTLZG3J2a/NVCAleEhjzq5oxgHyaCU9yYXvcLsvoVaHJq/s5xXI6/XXP6tz7R9xAOtHnSO/tXtF3WRTlA==", + "license": "MIT" + }, + "node_modules/asynckit": { + "version": "0.4.0", + "resolved": "https://registry.npmjs.org/asynckit/-/asynckit-0.4.0.tgz", + "integrity": "sha512-Oei9OH4tRh0YqU3GxhX79dM/mwVgvbZJaSNaRk+bshkj0S5cfHcgYakreBjrHwatXKbz+IoIdYLxrKim2MjW0Q==", + "license": "MIT" + }, + "node_modules/axios": { + "version": "1.16.1", + "resolved": "https://registry.npmjs.org/axios/-/axios-1.16.1.tgz", + "integrity": "sha512-caYkukvroVPO8KrzuJEb50Hm07KwfBZPEC3VeFHTsqWHvKTsy54hjJz9BS/cdaypROE2rH6xvm9mHX4fgWkr3A==", + "license": "MIT", + "dependencies": { + "follow-redirects": "^1.16.0", + "form-data": "^4.0.5", + "https-proxy-agent": "^5.0.1", + "proxy-from-env": "^2.1.0" + } + }, + "node_modules/call-bind-apply-helpers": { + "version": "1.0.2", + "resolved": "https://registry.npmjs.org/call-bind-apply-helpers/-/call-bind-apply-helpers-1.0.2.tgz", + "integrity": "sha512-Sp1ablJ0ivDkSzjcaJdxEunN5/XvksFJ2sMBFfq6x0ryhQV/2b/KwFe21cMpmHtPOSij8K99/wSfoEuTObmuMQ==", + "license": "MIT", + "dependencies": { + "es-errors": "^1.3.0", + "function-bind": "^1.1.2" + }, + "engines": { + "node": ">= 0.4" + } + }, + "node_modules/combined-stream": { + "version": "1.0.8", + "resolved": "https://registry.npmjs.org/combined-stream/-/combined-stream-1.0.8.tgz", + "integrity": "sha512-FQN4MRfuJeHf7cBbBMJFXhKSDq+2kAArBlmRBvcvFE5BB1HZKXtSFASDhdlz9zOYwxh8lDdnvmMOe/+5cdoEdg==", + "license": "MIT", + "dependencies": { + "delayed-stream": "~1.0.0" + }, + "engines": { + "node": ">= 0.8" + } + }, + "node_modules/cookie": { + "version": "0.7.2", + "resolved": "https://registry.npmjs.org/cookie/-/cookie-0.7.2.tgz", + "integrity": "sha512-yki5XnKuf750l50uGTllt6kKILY4nQ1eNIQatoXEByZ5dWgnKqbnqmTrBE5B4N7lrMJKQ2ytWMiTO2o0v6Ew/w==", + "license": "MIT", + "engines": { + "node": ">= 0.6" + } + }, + "node_modules/debug": { + "version": "2.6.9", + "resolved": "https://registry.npmjs.org/debug/-/debug-2.6.9.tgz", + "integrity": "sha512-bC7ElrdJaJnPbAP+1EotYvqZsb3ecl5wi6Bfi6BJTUcNowp6cvspg0jXznRTKDjm/E7AdgFBVeAPVMNcKGsHMA==", + "license": "MIT", + "dependencies": { + "ms": "2.0.0" + } + }, + "node_modules/delayed-stream": { + "version": "1.0.0", + "resolved": "https://registry.npmjs.org/delayed-stream/-/delayed-stream-1.0.0.tgz", + "integrity": "sha512-ZySD7Nf91aLB0RxL4KGrKHBXl7Eds1DAmEdcoVawXnLD7SDhpNgtuII2aAkg7a7QS41jxPSZ17p4VdGnMHk3MQ==", + "license": "MIT", + "engines": { + "node": ">=0.4.0" + } + }, + "node_modules/dunder-proto": { + "version": "1.0.1", + "resolved": "https://registry.npmjs.org/dunder-proto/-/dunder-proto-1.0.1.tgz", + "integrity": "sha512-KIN/nDJBQRcXw0MLVhZE9iQHmG68qAVIBg9CqmUYjmQIhgij9U5MFvrqkUL5FbtyyzZuOeOt0zdeRe4UY7ct+A==", + "license": "MIT", + "dependencies": { + "call-bind-apply-helpers": "^1.0.1", + "es-errors": "^1.3.0", + "gopd": "^1.2.0" + }, + "engines": { + "node": ">= 0.4" + } + }, + "node_modules/durable-functions": { + "version": "3.3.1", + "resolved": "https://registry.npmjs.org/durable-functions/-/durable-functions-3.3.1.tgz", + "integrity": "sha512-3LQ8ZE07p94ESQ9itgu1IdbXQ2qBeL/VTVnZV3VNyDqShJZWyr8nwCmhNLzp1xM5Uv+f//0iqJ6Gl26W2OGPsA==", + "license": "MIT", + "dependencies": { + "@azure/functions": "^4.0.0", + "@opentelemetry/api": "^1.9.0", + "axios": "^1.11.0", + "debug": "~2.6.9", + "lodash": "^4.17.15", + "moment": "^2.29.2", + "uuid": "^9.0.1", + "validator": "~13.15.20" + }, + "engines": { + "node": ">=18.0" + } + }, + "node_modules/es-define-property": { + "version": "1.0.1", + "resolved": "https://registry.npmjs.org/es-define-property/-/es-define-property-1.0.1.tgz", + "integrity": "sha512-e3nRfgfUZ4rNGL232gUgX06QNyyez04KdjFrF+LTRoOXmrOgFKDg4BCdsjW8EnT69eqdYGmRpJwiPVYNrCaW3g==", + "license": "MIT", + "engines": { + "node": ">= 0.4" + } + }, + "node_modules/es-errors": { + "version": "1.3.0", + "resolved": "https://registry.npmjs.org/es-errors/-/es-errors-1.3.0.tgz", + "integrity": "sha512-Zf5H2Kxt2xjTvbJvP2ZWLEICxA6j+hAmMzIlypy4xcBg1vKVnx89Wy0GbS+kf5cwCVFFzdCFh2XSCFNULS6csw==", + "license": "MIT", + "engines": { + "node": ">= 0.4" + } + }, + "node_modules/es-object-atoms": { + "version": "1.1.2", + "resolved": "https://registry.npmjs.org/es-object-atoms/-/es-object-atoms-1.1.2.tgz", + "integrity": "sha512-HWcBoN6NileqtSydK2FqHbS/LoDd2pqrnQHLyJzBj4kOp/ky2MWMN694xOfkK8/SnUsW2DH7EfyVlydKCsm1Zw==", + "license": "MIT", + "dependencies": { + "es-errors": "^1.3.0" + }, + "engines": { + "node": ">= 0.4" + } + }, + "node_modules/es-set-tostringtag": { + "version": "2.1.0", + "resolved": "https://registry.npmjs.org/es-set-tostringtag/-/es-set-tostringtag-2.1.0.tgz", + "integrity": "sha512-j6vWzfrGVfyXxge+O0x5sh6cvxAog0a/4Rdd2K36zCMV5eJ+/+tOAngRO8cODMNWbVRdVlmGZQL2YS3yR8bIUA==", + "license": "MIT", + "dependencies": { + "es-errors": "^1.3.0", + "get-intrinsic": "^1.2.6", + "has-tostringtag": "^1.0.2", + "hasown": "^2.0.2" + }, + "engines": { + "node": ">= 0.4" + } + }, + "node_modules/follow-redirects": { + "version": "1.16.0", + "resolved": "https://registry.npmjs.org/follow-redirects/-/follow-redirects-1.16.0.tgz", + "integrity": "sha512-y5rN/uOsadFT/JfYwhxRS5R7Qce+g3zG97+JrtFZlC9klX/W5hD7iiLzScI4nZqUS7DNUdhPgw4xI8W2LuXlUw==", + "funding": [ + { + "type": "individual", + "url": "https://github.com/sponsors/RubenVerborgh" + } + ], + "license": "MIT", + "engines": { + "node": ">=4.0" + }, + "peerDependenciesMeta": { + "debug": { + "optional": true + } + } + }, + "node_modules/form-data": { + "version": "4.0.5", + "resolved": "https://registry.npmjs.org/form-data/-/form-data-4.0.5.tgz", + "integrity": "sha512-8RipRLol37bNs2bhoV67fiTEvdTrbMUYcFTiy3+wuuOnUog2QBHCZWXDRijWQfAkhBj2Uf5UnVaiWwA5vdd82w==", + "license": "MIT", + "dependencies": { + "asynckit": "^0.4.0", + "combined-stream": "^1.0.8", + "es-set-tostringtag": "^2.1.0", + "hasown": "^2.0.2", + "mime-types": "^2.1.12" + }, + "engines": { + "node": ">= 6" + } + }, + "node_modules/function-bind": { + "version": "1.1.2", + "resolved": "https://registry.npmjs.org/function-bind/-/function-bind-1.1.2.tgz", + "integrity": "sha512-7XHNxH7qX9xG5mIwxkhumTox/MIRNcOgDrxWsMt2pAr23WHp6MrRlN7FBSFpCpr+oVO0F744iUgR82nJMfG2SA==", + "license": "MIT", + "funding": { + "url": "https://github.com/sponsors/ljharb" + } + }, + "node_modules/get-intrinsic": { + "version": "1.3.0", + "resolved": "https://registry.npmjs.org/get-intrinsic/-/get-intrinsic-1.3.0.tgz", + "integrity": "sha512-9fSjSaos/fRIVIp+xSJlE6lfwhES7LNtKaCBIamHsjr2na1BiABJPo0mOjjz8GJDURarmCPGqaiVg5mfjb98CQ==", + "license": "MIT", + "dependencies": { + "call-bind-apply-helpers": "^1.0.2", + "es-define-property": "^1.0.1", + "es-errors": "^1.3.0", + "es-object-atoms": "^1.1.1", + "function-bind": "^1.1.2", + "get-proto": "^1.0.1", + "gopd": "^1.2.0", + "has-symbols": "^1.1.0", + "hasown": "^2.0.2", + "math-intrinsics": "^1.1.0" + }, + "engines": { + "node": ">= 0.4" + }, + "funding": { + "url": "https://github.com/sponsors/ljharb" + } + }, + "node_modules/get-proto": { + "version": "1.0.1", + "resolved": "https://registry.npmjs.org/get-proto/-/get-proto-1.0.1.tgz", + "integrity": "sha512-sTSfBjoXBp89JvIKIefqw7U2CCebsc74kiY6awiGogKtoSGbgjYE/G/+l9sF3MWFPNc9IcoOC4ODfKHfxFmp0g==", + "license": "MIT", + "dependencies": { + "dunder-proto": "^1.0.1", + "es-object-atoms": "^1.0.0" + }, + "engines": { + "node": ">= 0.4" + } + }, + "node_modules/gopd": { + "version": "1.2.0", + "resolved": "https://registry.npmjs.org/gopd/-/gopd-1.2.0.tgz", + "integrity": "sha512-ZUKRh6/kUFoAiTAtTYPZJ3hw9wNxx+BIBOijnlG9PnrJsCcSjs1wyyD6vJpaYtgnzDrKYRSqf3OO6Rfa93xsRg==", + "license": "MIT", + "engines": { + "node": ">= 0.4" + }, + "funding": { + "url": "https://github.com/sponsors/ljharb" + } + }, + "node_modules/has-symbols": { + "version": "1.1.0", + "resolved": "https://registry.npmjs.org/has-symbols/-/has-symbols-1.1.0.tgz", + "integrity": "sha512-1cDNdwJ2Jaohmb3sg4OmKaMBwuC48sYni5HUw2DvsC8LjGTLK9h+eb1X6RyuOHe4hT0ULCW68iomhjUoKUqlPQ==", + "license": "MIT", + "engines": { + "node": ">= 0.4" + }, + "funding": { + "url": "https://github.com/sponsors/ljharb" + } + }, + "node_modules/has-tostringtag": { + "version": "1.0.2", + "resolved": "https://registry.npmjs.org/has-tostringtag/-/has-tostringtag-1.0.2.tgz", + "integrity": "sha512-NqADB8VjPFLM2V0VvHUewwwsw0ZWBaIdgo+ieHtK3hasLz4qeCRjYcqfB6AQrBggRKppKF8L52/VqdVsO47Dlw==", + "license": "MIT", + "dependencies": { + "has-symbols": "^1.0.3" + }, + "engines": { + "node": ">= 0.4" + }, + "funding": { + "url": "https://github.com/sponsors/ljharb" + } + }, + "node_modules/hasown": { + "version": "2.0.3", + "resolved": "https://registry.npmjs.org/hasown/-/hasown-2.0.3.tgz", + "integrity": "sha512-ej4AhfhfL2Q2zpMmLo7U1Uv9+PyhIZpgQLGT1F9miIGmiCJIoCgSmczFdrc97mWT4kVY72KA+WnnhJ5pghSvSg==", + "license": "MIT", + "dependencies": { + "function-bind": "^1.1.2" + }, + "engines": { + "node": ">= 0.4" + } + }, + "node_modules/https-proxy-agent": { + "version": "5.0.1", + "resolved": "https://registry.npmjs.org/https-proxy-agent/-/https-proxy-agent-5.0.1.tgz", + "integrity": "sha512-dFcAjpTQFgoLMzC2VwU+C/CbS7uRL0lWmxDITmqm7C+7F0Odmj6s9l6alZc6AELXhrnggM2CeWSXHGOdX2YtwA==", + "license": "MIT", + "dependencies": { + "agent-base": "6", + "debug": "4" + }, + "engines": { + "node": ">= 6" + } + }, + "node_modules/https-proxy-agent/node_modules/debug": { + "version": "4.4.3", + "resolved": "https://registry.npmjs.org/debug/-/debug-4.4.3.tgz", + "integrity": "sha512-RGwwWnwQvkVfavKVt22FGLw+xYSdzARwm0ru6DhTVA3umU5hZc28V3kO4stgYryrTlLpuvgI9GiijltAjNbcqA==", + "license": "MIT", + "dependencies": { + "ms": "^2.1.3" + }, + "engines": { + "node": ">=6.0" + }, + "peerDependenciesMeta": { + "supports-color": { + "optional": true + } + } + }, + "node_modules/https-proxy-agent/node_modules/ms": { + "version": "2.1.3", + "resolved": "https://registry.npmjs.org/ms/-/ms-2.1.3.tgz", + "integrity": "sha512-6FlzubTLZG3J2a/NVCAleEhjzq5oxgHyaCU9yYXvcLsvoVaHJq/s5xXI6/XXP6tz7R9xAOtHnSO/tXtF3WRTlA==", + "license": "MIT" + }, + "node_modules/lodash": { + "version": "4.18.1", + "resolved": "https://registry.npmjs.org/lodash/-/lodash-4.18.1.tgz", + "integrity": "sha512-dMInicTPVE8d1e5otfwmmjlxkZoUpiVLwyeTdUsi/Caj/gfzzblBcCE5sRHV/AsjuCmxWrte2TNGSYuCeCq+0Q==", + "license": "MIT" + }, + "node_modules/math-intrinsics": { + "version": "1.1.0", + "resolved": "https://registry.npmjs.org/math-intrinsics/-/math-intrinsics-1.1.0.tgz", + "integrity": "sha512-/IXtbwEk5HTPyEwyKX6hGkYXxM9nbj64B+ilVJnC/R6B0pH5G4V3b0pVbL7DBj4tkhBAppbQUlf6F6Xl9LHu1g==", + "license": "MIT", + "engines": { + "node": ">= 0.4" + } + }, + "node_modules/mime-db": { + "version": "1.52.0", + "resolved": "https://registry.npmjs.org/mime-db/-/mime-db-1.52.0.tgz", + "integrity": "sha512-sPU4uV7dYlvtWJxwwxHD0PuihVNiE7TyAbQ5SWxDCB9mUYvOgroQOwYQQOKPJ8CIbE+1ETVlOoK1UC2nU3gYvg==", + "license": "MIT", + "engines": { + "node": ">= 0.6" + } + }, + "node_modules/mime-types": { + "version": "2.1.35", + "resolved": "https://registry.npmjs.org/mime-types/-/mime-types-2.1.35.tgz", + "integrity": "sha512-ZDY+bPm5zTTF+YpCrAU9nK0UgICYPT0QtT1NZWFv4s++TNkcgVaT0g6+4R2uI4MjQjzysHB1zxuWL50hzaeXiw==", + "license": "MIT", + "dependencies": { + "mime-db": "1.52.0" + }, + "engines": { + "node": ">= 0.6" + } + }, + "node_modules/moment": { + "version": "2.30.1", + "resolved": "https://registry.npmjs.org/moment/-/moment-2.30.1.tgz", + "integrity": "sha512-uEmtNhbDOrWPFS+hdjFCBfy9f2YoyzRpwcl+DqpC6taX21FzsTLQVbMV/W7PzNSX6x/bhC1zA3c2UQ5NzH6how==", + "license": "MIT", + "engines": { + "node": "*" + } + }, + "node_modules/ms": { + "version": "2.0.0", + "resolved": "https://registry.npmjs.org/ms/-/ms-2.0.0.tgz", + "integrity": "sha512-Tpp60P6IUJDTuOq/5Z8cdskzJujfwqfOTkrwIwj7IRISpnkJnT6SyJ4PCPnGMoFjC9ddhal5KVIYtAt97ix05A==", + "license": "MIT" + }, + "node_modules/proxy-from-env": { + "version": "2.1.0", + "resolved": "https://registry.npmjs.org/proxy-from-env/-/proxy-from-env-2.1.0.tgz", + "integrity": "sha512-cJ+oHTW1VAEa8cJslgmUZrc+sjRKgAKl3Zyse6+PV38hZe/V6Z14TbCuXcan9F9ghlz4QrFr2c92TNF82UkYHA==", + "license": "MIT", + "engines": { + "node": ">=10" + } + }, + "node_modules/uuid": { + "version": "9.0.1", + "resolved": "https://registry.npmjs.org/uuid/-/uuid-9.0.1.tgz", + "integrity": "sha512-b+1eJOlsR9K8HJpow9Ok3fiWOWSIcIzXodvv0rQjVoOVNpWMpxf1wZNpt4y9h10odCNrqnYp1OBzRktckBe3sA==", + "deprecated": "uuid@10 and below is no longer supported. For ESM codebases, update to uuid@latest. For CommonJS codebases, use uuid@11 (but be aware this version will likely be deprecated in 2028).", + "funding": [ + "https://github.com/sponsors/broofa", + "https://github.com/sponsors/ctavan" + ], + "license": "MIT", + "bin": { + "uuid": "dist/bin/uuid" + } + }, + "node_modules/validator": { + "version": "13.15.35", + "resolved": "https://registry.npmjs.org/validator/-/validator-13.15.35.tgz", + "integrity": "sha512-TQ5pAGhd5whStmqWvYF4OjQROlmv9SMFVt37qoCBdqRffuuklWYQlCNnEs2ZaIBD1kZRNnikiZOS1eqgkar0iw==", + "license": "MIT", + "engines": { + "node": ">= 0.10" + } + } + } +} diff --git a/samples/scenarios/WorkItemFilteringSplitActivitiesPython/.gitignore b/samples/scenarios/WorkItemFilteringSplitActivitiesPython/.gitignore new file mode 100644 index 00000000..1f0dd2de --- /dev/null +++ b/samples/scenarios/WorkItemFilteringSplitActivitiesPython/.gitignore @@ -0,0 +1,21 @@ +# Byte-compiled / optimized / DLL files +__pycache__/ +*.py[cod] +*$py.class + +# Virtual environments +.venv/ +venv/ +env/ +ENV/ + +# Distribution / packaging +build/ +dist/ +*.egg-info/ + +# Local run logs +.logs/ + +# Azure Developer CLI +.azure diff --git a/samples/scenarios/WorkItemFilteringSplitActivitiesPython/README.md b/samples/scenarios/WorkItemFilteringSplitActivitiesPython/README.md new file mode 100644 index 00000000..fc18935d --- /dev/null +++ b/samples/scenarios/WorkItemFilteringSplitActivitiesPython/README.md @@ -0,0 +1,316 @@ +# Work Item Filtering — Split Activities Sample (Python) + +This sample demonstrates **Work Item Filtering**, a feature that allows workers to declare which orchestrations, activities, and entities they can process. The Durable Task Scheduler (DTS) backend routes work items only to workers whose filters match, preventing workers from receiving work they cannot handle. + +Before work item filtering, all orchestrations, activities, and entities were handed to any connected worker regardless of what it actually hosted. This caused errors (or silent hangs) when a worker received a work item it didn't implement — especially problematic in multi-service deployments, rolling upgrades, and microservice topologies. With filtering, each worker registers its task set; DTS creates per-filter queues and routes work items to matching workers. If no filter is specified, a worker is eligible to receive any work item type or name, and DTS dispatches each work item to one eligible worker (it is not broadcast to every connected worker). + +This is the Python version of the sample. Equivalent [.NET](../WorkItemFilteringSplitActivities/) and [Java](../WorkItemFilteringSplitActivitiesJava/) versions are also available. + +## Architecture + +``` +┌─────────────────────────────────────────────────────────────┐ +│ Durable Task Scheduler (DTS) │ +│ │ +│ Orchestration queue ──► routed to Orchestrator Worker only │ +│ validate_order queue ─► routed to Validator Worker only │ +│ ship_order queue ──► routed to Shipper Worker only │ +└────────────┬──────────────────┬──────────────────┬──────────┘ + │ │ │ + ┌───────▼───────┐ ┌──────▼───────┐ ┌──────▼───────┐ + │ Orchestrator │ │ Validator │ │ Shipper │ + │ Worker │ │ Worker │ │ Worker │ + │ │ │ │ │ │ + │ Registers: │ │ Registers: │ │ Registers: │ + │ • order_proc- │ │ • validate_ │ │ • ship_order │ + │ essing_ │ │ order │ │ │ + │ orchestrator │ │ │ │ │ + └───────────────┘ └───────────────┘ └───────────────┘ + + ┌───────────────┐ + │ Client │ + │ (Driver) │ + │ │ + │ Schedules new │ + │ orchestrations │ + │ and prints │ + │ results │ + └───────────────┘ +``` + +**Orchestrator Worker** runs orchestrations only — it has no activities registered. +**Validator Worker** runs `validate_order` only — it has no orchestrations or other activities. +**Shipper Worker** runs `ship_order` only — same isolation. +**Client** schedules orchestrations and polls for completion. + +## The Orchestration + +`order_processing_orchestrator` performs two sequential activity calls: + +1. `validate_order(order_id)` → routed to Validator Worker +2. `ship_order(order_id)` → routed to Shipper Worker + +Returns a combined result string. The orchestrator calls each activity **by name**, so it doesn't need to import or register the activity functions — they run in separate worker processes. + +## Prerequisites + +- [Python 3.10+](https://www.python.org/downloads/) +- [Docker](https://docs.docker.com/get-docker/) (for the DTS emulator) +- [Azure Developer CLI (`azd`)](https://learn.microsoft.com/azure/developer/azure-developer-cli/) (for deploying to Azure) + +## Running Locally + +### Option A: One command + +```bash +cd samples/scenarios/WorkItemFilteringSplitActivitiesPython +./run-local.sh +``` + +This starts the DTS emulator, creates a virtual environment, installs dependencies, and launches all three workers plus the client, tailing their logs. Press Ctrl+C to stop everything. + +### Option B: Manual steps + +#### 1. Start the DTS Emulator + +```bash +docker pull mcr.microsoft.com/dts/dts-emulator:latest +docker run -d --name dts-emulator -p 8080:8080 -p 8082:8082 mcr.microsoft.com/dts/dts-emulator:latest +``` + +The emulator dashboard is available at `http://localhost:8082`. + +#### 2. Create a virtual environment and install dependencies + +```bash +cd samples/scenarios/WorkItemFilteringSplitActivitiesPython +python -m venv .venv +source .venv/bin/activate # On Windows: .venv\Scripts\activate +pip install -r src/client/requirements.txt +``` + +All four services share the same two dependencies (`durabletask-azuremanaged` and `azure-identity`), so a single install covers everything. + +#### 3. Start the three workers (each in a separate terminal) + +**Terminal 1 — Orchestrator Worker:** +```bash +python src/orchestrator-worker/orchestrator_worker.py +``` + +**Terminal 2 — Validator Worker (validate_order activity):** +```bash +python src/validator-worker/validator_worker.py +``` + +**Terminal 3 — Shipper Worker (ship_order activity):** +```bash +python src/shipper-worker/shipper_worker.py +``` + +#### 4. Run the Client (in a fourth terminal) + +```bash +python src/client/client.py +``` + +## Expected Output + +The client runs in a **continuous loop**, scheduling a batch of 3 orchestrations every 30 seconds for 10 minutes. This makes it easy to observe scaling behavior over time. + +### Client terminal + +``` +10:30:01 === Work Item Filtering Demo — Client === +10:30:01 Will schedule 3 orchestrations every 30s for 10 minutes. + +10:30:01 --- Batch #1 at 10:30:01 --- +10:30:01 Scheduling orchestration with orderId='ORD-B001-001'... +10:30:01 -> Scheduled with InstanceId=abc123 +10:30:01 Scheduling orchestration with orderId='ORD-B001-002'... +10:30:01 -> Scheduled with InstanceId=def456 +10:30:01 Scheduling orchestration with orderId='ORD-B001-003'... +10:30:01 -> Scheduled with InstanceId=ghi789 +10:30:02 COMPLETED | InstanceId=abc123 | Output: "Order 'ORD-B001-001' => Validation: [Order ORD-B001-001 is valid], Shipping: [Shipped with tracking TRACK-ORD-B001-001-4271]" +10:30:02 Batch #1 results: 3 completed, 0 failed +10:30:02 Next batch in 30s (deadline in 10.0 min) +``` + +### Orchestrator Worker terminal (orchestrations only — no activities) + +``` +10:30:02 [Orchestrator] Orchestration | Name=order_processing_orchestrator | InstanceId=abc123 | Processing order 'ORD-B001-001' +10:30:02 [Orchestrator] Orchestration | InstanceId=abc123 | Dispatching validate_order to Validator Worker... +10:30:02 [Orchestrator] Orchestration | InstanceId=abc123 | Dispatching ship_order to Shipper Worker... +10:30:02 [Orchestrator] Orchestration | InstanceId=abc123 | Completed: Order 'ORD-B001-001' => Validation: [...], Shipping: [...] +``` + +### Validator Worker terminal (validate_order only — no ship_order, no orchestrations) + +``` +10:30:02 [Validator] Activity | Name=validate_order | InstanceId=abc123 | Validating order 'ORD-B001-001'... +10:30:02 [Validator] Activity | Name=validate_order | InstanceId=abc123 | Result: Order ORD-B001-001 is valid +``` + +### Shipper Worker terminal (ship_order only — no validate_order, no orchestrations) + +``` +10:30:02 [Shipper] Activity | Name=ship_order | InstanceId=abc123 | Shipping order 'ORD-B001-001'... +10:30:02 [Shipper] Activity | Name=ship_order | InstanceId=abc123 | Result: Shipped with tracking TRACK-ORD-B001-001-4271 +``` + +**Key observation:** Each worker processes **only** its registered work item types. No cross-processing occurs. + +## What to Try Next: Strict Routing Experiment + +1. **Stop Shipper Worker** (Ctrl+C in Terminal 3). +2. Let the Client schedule new orchestrations. +3. Observe that: + - Orchestrator Worker picks up and starts orchestrations. + - Validator Worker completes `validate_order` for each order. + - `ship_order` work items **remain pending** — they are not delivered to Validator Worker or Orchestrator Worker. + - The orchestrations stay in "Running" status, waiting for the `ship_order` activity to complete. +4. **Restart Shipper Worker** — the pending `ship_order` work items are immediately delivered and the orchestrations complete. + +This demonstrates that filtering is strict: work items are routed only to workers with matching filters. There is no fallback to other workers. + +## How It Works + +Each worker process registers its tasks with the `DurableTaskSchedulerWorker` and then calls `use_work_item_filters()`. The SDK automatically constructs **work item filters** from whatever is registered: + +- Orchestrator Worker's filter: `orchestrations: [order_processing_orchestrator]` +- Validator Worker's filter: `activities: [validate_order]` +- Shipper Worker's filter: `activities: [ship_order]` + +```python +# Orchestrator Worker — registers only the orchestrator +worker.add_orchestrator(order_processing_orchestrator) +worker.use_work_item_filters() # Auto-generate filters from the registry +``` + +DTS creates per-filter queues and routes each work item to the matching queue. If a filter list is empty for a given type (e.g., Validator Worker has no orchestration filter), that worker simply never receives work items of that type. + +For more control, you can pass explicit `WorkItemFilters` with optional version constraints: + +```python +from durabletask import worker + +worker.use_work_item_filters(worker.WorkItemFilters( + activities=[ + worker.ActivityWorkItemFilter(name="validate_order"), + ], +)) +``` + +## Deploying to Azure + +This sample includes full infrastructure-as-code (Bicep) and an `azure.yaml` for one-command deployment via [Azure Developer CLI (`azd`)](https://learn.microsoft.com/azure/developer/azure-developer-cli/). + +### What Gets Deployed + +| Resource | Purpose | +|---|---| +| **Resource Group** | Contains all resources | +| **Durable Task Scheduler** (Consumption SKU) | Managed orchestration backend | +| **Task Hub** | Logical unit for orchestrations and work items | +| **Container Apps Environment** | Shared hosting environment with VNet integration | +| **Azure Container Registry** | Stores Docker images for each service | +| **User-Assigned Managed Identity** | Shared identity with DTS Worker/Client RBAC role | +| **4 Container Apps** | Client, Orchestrator Worker, Validator Worker, Shipper Worker | + +### Deploy with `azd` + +```bash +cd samples/scenarios/WorkItemFilteringSplitActivitiesPython +azd up +``` + +You'll be prompted for an environment name, subscription, and location. The deployment takes ~5 minutes. + +### KEDA Scaling with DTS + +Each worker Container App is configured with a **DTS-aware KEDA custom scale rule** (`azure-durabletask-scheduler`) that scales based on the **work item backlog** in the task hub. The key parameter is `workItemType`, which tells the scaler what kind of work to monitor: + +| Container App | Service Name | `workItemType` | Scales on | +|---|---|---|---| +| **Client** | `client` | `Orchestration` | Pending orchestration work items | +| **Orchestrator Worker** | `orchestrator-worker` | `Orchestration` | Pending orchestration work items | +| **Validator Worker** | `validator-worker` | `Activity` | Pending activity work items | +| **Shipper Worker** | `shipper-worker` | `Activity` | Pending activity work items | + +The scale rule metadata (from [app.bicep](infra/app/app.bicep)): + +```bicep +scaleRuleType: 'azure-durabletask-scheduler' +scaleRuleMetadata: { + endpoint: dtsEndpoint // DTS scheduler URL + maxConcurrentWorkItemsCount: '1' + taskhubName: taskHubName + workItemType: workItemType // 'Orchestration' or 'Activity' +} +scaleRuleIdentity: userAssignedManagedIdentity.resourceId +``` + +- Workers scale from **0 to 10** replicas. When the client finishes its loop and no more work items arrive, workers scale back to zero. +- The `scaleRuleIdentity` uses the shared user-assigned managed identity to authenticate with DTS, so no connection strings or secrets are needed for scaling. +- `maxConcurrentWorkItemsCount: '1'` means KEDA will scale up one replica per pending work item, up to the max. + +### Manual Deployment (without `azd`) + +Set the `ENDPOINT` and `TASKHUB` environment variables to point to your deployed scheduler: + +```bash +export ENDPOINT="https://your-scheduler.westus2.durabletask.io" +export TASKHUB="your-taskhub-name" +``` + +The workers and client will automatically use `DefaultAzureCredential` for authentication. Make sure the identity running each process has the **Durable Task Scheduler Worker** / **Durable Task Scheduler Client** role on the scheduler resource. When running inside Container Apps, the `AZURE_MANAGED_IDENTITY_CLIENT_ID` environment variable selects the shared user-assigned managed identity. + +## Project Structure + +``` +WorkItemFilteringSplitActivitiesPython/ +├── README.md +├── azure.yaml # azd service definitions +├── run-local.sh # Local run helper (emulator + all services) +├── .gitignore +├── infra/ # Bicep infrastructure-as-code +│ ├── main.bicep # Top-level — resource group, DTS, container apps +│ ├── main.parameters.json +│ ├── abbreviations.json +│ ├── app/ +│ │ ├── app.bicep # Per-service container app (with KEDA scale rule) +│ │ ├── dts.bicep # DTS scheduler + task hub +│ │ └── user-assigned-identity.bicep +│ └── core/ +│ ├── host/ # Container Apps Environment, Registry, App template +│ ├── networking/ # VNet +│ └── security/ # ACR pull role, DTS role assignments +└── src/ + ├── client/ # Schedules orchestrations in a loop, prints results + │ ├── client.py + │ ├── requirements.txt + │ └── Dockerfile + ├── orchestrator-worker/ # Orchestrator Worker — runs orchestrations only + │ ├── orchestrator_worker.py + │ ├── requirements.txt + │ └── Dockerfile + ├── validator-worker/ # Validator Worker — runs validate_order activity only + │ ├── validator_worker.py + │ ├── requirements.txt + │ └── Dockerfile + └── shipper-worker/ # Shipper Worker — runs ship_order activity only + ├── shipper_worker.py + ├── requirements.txt + └── Dockerfile +``` + +## Viewing in the Dashboard + +- **Emulator:** Navigate to `http://localhost:8082` → select the "default" task hub. +- **Azure:** Navigate to your Scheduler resource in the Azure Portal → Task Hub → Dashboard URL. + +## Reference + +- [Durable Task Scheduler documentation](https://learn.microsoft.com/azure/azure-functions/durable/durable-task-scheduler/develop-with-durable-task-scheduler) +- [Durable Task Python SDK](https://github.com/microsoft/durabletask-python) diff --git a/samples/scenarios/WorkItemFilteringSplitActivitiesPython/azure.yaml b/samples/scenarios/WorkItemFilteringSplitActivitiesPython/azure.yaml new file mode 100644 index 00000000..1c4451cc --- /dev/null +++ b/samples/scenarios/WorkItemFilteringSplitActivitiesPython/azure.yaml @@ -0,0 +1,34 @@ +# yaml-language-server: $schema=https://raw.githubusercontent.com/Azure/azure-dev/main/schemas/v1.0/azure.yaml.json + +metadata: + template: work-item-filtering-split-activities-python +name: work-item-filtering-split-activities-python +services: + client: + project: ./src/client + language: python + host: containerapp + apiVersion: 2025-01-01 + docker: + path: ./Dockerfile + orchestrator-worker: + project: ./src/orchestrator-worker + language: python + host: containerapp + apiVersion: 2025-01-01 + docker: + path: ./Dockerfile + validator-worker: + project: ./src/validator-worker + language: python + host: containerapp + apiVersion: 2025-01-01 + docker: + path: ./Dockerfile + shipper-worker: + project: ./src/shipper-worker + language: python + host: containerapp + apiVersion: 2025-01-01 + docker: + path: ./Dockerfile diff --git a/samples/scenarios/WorkItemFilteringSplitActivitiesPython/infra/abbreviations.json b/samples/scenarios/WorkItemFilteringSplitActivitiesPython/infra/abbreviations.json new file mode 100644 index 00000000..1f9a112f --- /dev/null +++ b/samples/scenarios/WorkItemFilteringSplitActivitiesPython/infra/abbreviations.json @@ -0,0 +1,139 @@ +{ + "analysisServicesServers": "as", + "apiManagementService": "apim-", + "appConfigurationStores": "appcs-", + "appManagedEnvironments": "cae-", + "appContainerApps": "ca-", + "authorizationPolicyDefinitions": "policy-", + "automationAutomationAccounts": "aa-", + "blueprintBlueprints": "bp-", + "blueprintBlueprintsArtifacts": "bpa-", + "cacheRedis": "redis-", + "cdnProfiles": "cdnp-", + "cdnProfilesEndpoints": "cdne-", + "cognitiveServicesAccounts": "cog-", + "cognitiveServicesFormRecognizer": "cog-fr-", + "cognitiveServicesTextAnalytics": "cog-ta-", + "cognitiveServicesSpeech": "cog-sp-", + "computeAvailabilitySets": "avail-", + "computeCloudServices": "cld-", + "computeDiskEncryptionSets": "des", + "computeDisks": "disk", + "computeDisksOs": "osdisk", + "computeGalleries": "gal", + "computeSnapshots": "snap-", + "computeVirtualMachines": "vm", + "computeVirtualMachineScaleSets": "vmss-", + "containerInstanceContainerGroups": "ci", + "containerRegistryRegistries": "cr", + "containerServiceManagedClusters": "aks-", + "databricksWorkspaces": "dbw-", + "dataFactoryFactories": "adf-", + "dataLakeAnalyticsAccounts": "dla", + "dataLakeStoreAccounts": "dls", + "dataMigrationServices": "dms-", + "dBforMySQLServers": "mysql-", + "dBforPostgreSQLServers": "psql-", + "devicesIotHubs": "iot-", + "devicesProvisioningServices": "provs-", + "devicesProvisioningServicesCertificates": "pcert-", + "documentDBDatabaseAccounts": "cosmos-", + "eventGridDomains": "evgd-", + "eventGridDomainsTopics": "evgt-", + "eventGridEventSubscriptions": "evgs-", + "eventHubNamespaces": "evhns-", + "eventHubNamespacesEventHubs": "evh-", + "hdInsightClustersHadoop": "hadoop-", + "hdInsightClustersHbase": "hbase-", + "hdInsightClustersKafka": "kafka-", + "hdInsightClustersMl": "mls-", + "hdInsightClustersSpark": "spark-", + "hdInsightClustersStorm": "storm-", + "hybridComputeMachines": "arcs-", + "insightsActionGroups": "ag-", + "insightsComponents": "appi-", + "keyVaultVaults": "kv-", + "kubernetesConnectedClusters": "arck", + "kustoClusters": "dec", + "kustoClustersDatabases": "dedb", + "loadTesting": "lt-", + "logicIntegrationAccounts": "ia-", + "logicWorkflows": "logic-", + "machineLearningServicesWorkspaces": "mlw-", + "managedIdentityUserAssignedIdentities": "id-", + "managementManagementGroups": "mg-", + "migrateAssessmentProjects": "migr-", + "networkApplicationGateways": "agw-", + "networkApplicationSecurityGroups": "asg-", + "networkAzureFirewalls": "afw-", + "networkBastionHosts": "bas-", + "networkConnections": "con-", + "networkDnsZones": "dnsz-", + "networkExpressRouteCircuits": "erc-", + "networkFirewallPolicies": "afwp-", + "networkFirewallPoliciesWebApplication": "waf", + "networkFirewallPoliciesRuleGroups": "wafrg", + "networkFrontDoors": "fd-", + "networkFrontdoorWebApplicationFirewallPolicies": "fdfp-", + "networkLoadBalancersExternal": "lbe-", + "networkLoadBalancersInternal": "lbi-", + "networkLoadBalancersInboundNatRules": "rule-", + "networkLocalNetworkGateways": "lgw-", + "networkNatGateways": "ng-", + "networkNetworkInterfaces": "nic-", + "networkNetworkSecurityGroups": "nsg-", + "networkNetworkSecurityGroupsSecurityRules": "nsgsr-", + "networkNetworkWatchers": "nw-", + "networkPrivateDnsZones": "pdnsz-", + "networkPrivateLinkServices": "pl-", + "networkPublicIPAddresses": "pip-", + "networkPublicIPPrefixes": "ippre-", + "networkRouteFilters": "rf-", + "networkRouteTables": "rt-", + "networkRouteTablesRoutes": "udr-", + "networkTrafficManagerProfiles": "traf-", + "networkVirtualNetworkGateways": "vgw-", + "networkVirtualNetworks": "vnet-", + "networkVirtualNetworksSubnets": "snet-", + "networkVirtualNetworksVirtualNetworkPeerings": "peer-", + "networkVirtualWans": "vwan-", + "networkVpnGateways": "vpng-", + "networkVpnGatewaysVpnConnections": "vcn-", + "networkVpnGatewaysVpnSites": "vst-", + "notificationHubsNamespaces": "ntfns-", + "notificationHubsNamespacesNotificationHubs": "ntf-", + "operationalInsightsWorkspaces": "log-", + "portalDashboards": "dash-", + "powerBIDedicatedCapacities": "pbi-", + "purviewAccounts": "pview-", + "recoveryServicesVaults": "rsv-", + "resourcesResourceGroups": "rg-", + "searchSearchServices": "srch-", + "serviceBusNamespaces": "sb-", + "serviceBusNamespacesQueues": "sbq-", + "serviceBusNamespacesTopics": "sbt-", + "serviceEndPointPolicies": "se-", + "serviceFabricClusters": "sf-", + "signalRServiceSignalR": "sigr", + "sqlManagedInstances": "sqlmi-", + "sqlServers": "sql-", + "sqlServersDataWarehouse": "sqldw-", + "sqlServersDatabases": "sqldb-", + "sqlServersDatabasesStretch": "sqlstrdb-", + "storageStorageAccounts": "st", + "storageStorageAccountsVm": "stvm", + "storSimpleManagers": "ssimp", + "streamAnalyticsCluster": "asa-", + "synapseWorkspaces": "syn", + "synapseWorkspacesAnalyticsWorkspaces": "synw", + "synapseWorkspacesSqlPoolsDedicated": "syndp", + "synapseWorkspacesSqlPoolsSpark": "synsp", + "timeSeriesInsightsEnvironments": "tsi-", + "webServerFarms": "plan-", + "webSitesAppService": "app-", + "webSitesAppServiceEnvironment": "ase-", + "webSitesFunctions": "func-", + "webStaticSites": "stapp-", + "dts": "dts-", + "taskhub": "taskhub-" +} diff --git a/samples/scenarios/WorkItemFilteringSplitActivitiesPython/infra/app/app.bicep b/samples/scenarios/WorkItemFilteringSplitActivitiesPython/infra/app/app.bicep new file mode 100644 index 00000000..ec4953a4 --- /dev/null +++ b/samples/scenarios/WorkItemFilteringSplitActivitiesPython/infra/app/app.bicep @@ -0,0 +1,66 @@ +param appName string +param location string = resourceGroup().location +param tags object = {} + +param identityName string +param containerAppsEnvironmentName string +param containerRegistryName string +param serviceName string = 'aca' +param dtsEndpoint string +param taskHubName string + +@description('DTS work item type for KEDA scaling: Orchestration or Activity') +param workItemType string = 'Orchestration' + +type managedIdentity = { + resourceId: string + clientId: string +} + +@description('Unique identifier for user-assigned managed identity.') +param userAssignedManagedIdentity managedIdentity + +module containerAppsApp '../core/host/container-app.bicep' = { + name: 'container-apps-${serviceName}' + params: { + name: appName + containerAppsEnvironmentName: containerAppsEnvironmentName + containerRegistryName: containerRegistryName + location: location + tags: union(tags, { 'azd-service-name': serviceName }) + ingressEnabled: false + secrets: { + 'azure-managed-identity-client-id': userAssignedManagedIdentity.clientId + } + env: [ + { + name: 'AZURE_MANAGED_IDENTITY_CLIENT_ID' + secretRef: 'azure-managed-identity-client-id' + } + { + name: 'ENDPOINT' + value: dtsEndpoint + } + { + name: 'TASKHUB' + value: taskHubName + } + ] + identityName: identityName + containerMinReplicas: 0 + containerMaxReplicas: 10 + enableCustomScaleRule: true + scaleRuleName: 'dtsscaler-${serviceName}' + scaleRuleType: 'azure-durabletask-scheduler' + scaleRuleMetadata: { + endpoint: dtsEndpoint + maxConcurrentWorkItemsCount: '1' + taskhubName: taskHubName + workItemType: workItemType + } + scaleRuleIdentity: userAssignedManagedIdentity.resourceId + } +} + +output endpoint string = containerAppsApp.outputs.uri +output envName string = containerAppsApp.outputs.name diff --git a/samples/scenarios/WorkItemFilteringSplitActivitiesPython/infra/app/dts.bicep b/samples/scenarios/WorkItemFilteringSplitActivitiesPython/infra/app/dts.bicep new file mode 100644 index 00000000..6280e5a4 --- /dev/null +++ b/samples/scenarios/WorkItemFilteringSplitActivitiesPython/infra/app/dts.bicep @@ -0,0 +1,31 @@ +param ipAllowlist array +param location string +param tags object = {} +param name string +param taskhubname string +param skuName string +param skuCapacity int = 0 + +resource dts 'Microsoft.DurableTask/schedulers@2025-11-01' = { + location: location + tags: tags + name: name + properties: { + ipAllowlist: ipAllowlist + sku: skuName == 'Dedicated' ? { + name: skuName + capacity: skuCapacity + } : { + name: skuName + } + } +} + +resource taskhub 'Microsoft.DurableTask/schedulers/taskhubs@2025-11-01' = { + parent: dts + name: taskhubname +} + +output dts_NAME string = dts.name +output dts_URL string = dts.properties.endpoint +output TASKHUB_NAME string = taskhub.name diff --git a/samples/scenarios/WorkItemFilteringSplitActivitiesPython/infra/app/user-assigned-identity.bicep b/samples/scenarios/WorkItemFilteringSplitActivitiesPython/infra/app/user-assigned-identity.bicep new file mode 100644 index 00000000..0583ab8d --- /dev/null +++ b/samples/scenarios/WorkItemFilteringSplitActivitiesPython/infra/app/user-assigned-identity.bicep @@ -0,0 +1,17 @@ +metadata description = 'Creates a Microsoft Entra user-assigned identity.' + +param name string +param location string = resourceGroup().location +param tags object = {} + +resource identity 'Microsoft.ManagedIdentity/userAssignedIdentities@2023-01-31' = { + name: name + location: location + tags: tags +} + +output name string = identity.name +output resourceId string = identity.id +output principalId string = identity.properties.principalId +output clientId string = identity.properties.clientId +output tenantId string = identity.properties.tenantId diff --git a/samples/scenarios/WorkItemFilteringSplitActivitiesPython/infra/core/host/container-app.bicep b/samples/scenarios/WorkItemFilteringSplitActivitiesPython/infra/core/host/container-app.bicep new file mode 100644 index 00000000..8b1df5ca --- /dev/null +++ b/samples/scenarios/WorkItemFilteringSplitActivitiesPython/infra/core/host/container-app.bicep @@ -0,0 +1,195 @@ +metadata description = 'Creates a container app in an Azure Container App environment.' +param name string +param location string = resourceGroup().location +param tags object = {} + +@description('Allowed origins') +param allowedOrigins array = [] + +@description('Name of the environment for container apps') +param containerAppsEnvironmentName string + +@description('CPU cores allocated to a single container instance, e.g., 0.5') +param containerCpuCoreCount string = '0.5' + +@description('The maximum number of replicas to run. Must be at least 1.') +@minValue(1) +param containerMaxReplicas int = 10 + +@description('Memory allocated to a single container instance, e.g., 1Gi') +param containerMemory string = '1.0Gi' + +@description('The minimum number of replicas to run. Must be at least 0.') +param containerMinReplicas int = 0 + +@description('The name of the container') +param containerName string = 'main' + +@description('Enable custom scale rule') +param enableCustomScaleRule bool = false + +@description('Scale rule name') +param scaleRuleName string = 'scaler' + +@description('Scale rule type') +param scaleRuleType string = '' + +@description('Scale rule metadata') +param scaleRuleMetadata object = {} + +@description('Scale rule identity') +param scaleRuleIdentity string = '' + +@description('The name of the container registry') +param containerRegistryName string = '' + +@description('Hostname suffix for container registry. Set when deploying to sovereign clouds') +param containerRegistryHostSuffix string = 'azurecr.io' + +@description('The protocol used by Dapr to connect to the app, e.g., http or grpc') +@allowed([ 'http', 'grpc' ]) +param daprAppProtocol string = 'http' + +@description('The Dapr app ID') +param daprAppId string = containerName + +@description('Enable Dapr') +param daprEnabled bool = false + +@description('The environment variables for the container') +param env array = [] + +@description('Specifies if the resource ingress is exposed externally') +param external bool = true + +@description('The name of the user-assigned identity') +param identityName string = '' + +@description('The type of identity for the resource') +@allowed([ 'None', 'SystemAssigned', 'UserAssigned' ]) +param identityType string = 'None' + +@description('The name of the container image') +param imageName string = '' + +@description('Specifies if Ingress is enabled for the container app') +param ingressEnabled bool = true + +param revisionMode string = 'Single' + +@description('The secrets required for the container') +@secure() +param secrets object = {} + +@description('The service binds associated with the container') +param serviceBinds array = [] + +@description('The name of the container apps add-on to use. e.g. redis') +param serviceType string = '' + +@description('The target port for the container') +param targetPort int = 80 + +resource userIdentity 'Microsoft.ManagedIdentity/userAssignedIdentities@2023-01-31' existing = if (!empty(identityName)) { + name: identityName +} + +// Private registry support requires both an ACR name and a User Assigned managed identity +var usePrivateRegistry = !empty(identityName) && !empty(containerRegistryName) + +// Automatically set to `UserAssigned` when an `identityName` has been set +var normalizedIdentityType = !empty(identityName) ? 'UserAssigned' : identityType + +module containerRegistryAccess '../security/registry-access.bicep' = if (usePrivateRegistry) { + name: '${deployment().name}-registry-access' + params: { + containerRegistryName: containerRegistryName + principalId: usePrivateRegistry ? userIdentity!.properties.principalId : '' + } +} + +resource app 'Microsoft.App/containerApps@2025-01-01' = { + name: name + location: location + tags: tags + // It is critical that the identity is granted ACR pull access before the app is created + // otherwise the container app will throw a provision error + // This also forces us to use an user assigned managed identity since there would no way to + // provide the system assigned identity with the ACR pull access before the app is created + dependsOn: usePrivateRegistry ? [ containerRegistryAccess ] : [] + identity: { + type: normalizedIdentityType + userAssignedIdentities: !empty(identityName) && normalizedIdentityType == 'UserAssigned' ? { '${userIdentity.id}': {} } : null + } + properties: { + managedEnvironmentId: containerAppsEnvironment.id + configuration: { + activeRevisionsMode: revisionMode + ingress: ingressEnabled ? { + external: external + targetPort: targetPort + transport: 'auto' + corsPolicy: { + allowedOrigins: union([ 'https://portal.azure.com', 'https://ms.portal.azure.com' ], allowedOrigins) + } + } : null + dapr: daprEnabled ? { + enabled: true + appId: daprAppId + appProtocol: daprAppProtocol + appPort: ingressEnabled ? targetPort : 0 + } : { enabled: false } + secrets: [for secret in items(secrets): { + name: secret.key + #disable-next-line use-secure-value-for-secure-inputs + value: secret.value + }] + service: !empty(serviceType) ? { type: serviceType } : null + registries: usePrivateRegistry ? [ + { + server: '${containerRegistryName}.${containerRegistryHostSuffix}' + identity: userIdentity.id + } + ] : [] + } + template: { + serviceBinds: !empty(serviceBinds) ? serviceBinds : null + containers: [ + { + image: !empty(imageName) ? imageName : 'mcr.microsoft.com/azuredocs/containerapps-helloworld:latest' + name: containerName + env: env + resources: { + cpu: json(containerCpuCoreCount) + memory: containerMemory + } + } + ] + scale: { + minReplicas: containerMinReplicas + maxReplicas: containerMaxReplicas + rules: enableCustomScaleRule ? [ + { + name: scaleRuleName + custom: { + type: scaleRuleType + metadata: scaleRuleMetadata + identity: scaleRuleIdentity + } + } + ] : [] + } + } + } +} + +resource containerAppsEnvironment 'Microsoft.App/managedEnvironments@2023-05-01' existing = { + name: containerAppsEnvironmentName +} + +output defaultDomain string = containerAppsEnvironment.properties.defaultDomain +output identityPrincipalId string = normalizedIdentityType == 'None' ? '' : (empty(identityName) ? app.identity.principalId : userIdentity!.properties.principalId) +output imageName string = imageName +output name string = app.name +output serviceBind object = !empty(serviceType) ? { serviceId: app.id, name: name } : {} +output uri string = ingressEnabled ? 'https://${app.properties.configuration.ingress.fqdn}' : '' diff --git a/samples/scenarios/WorkItemFilteringSplitActivitiesPython/infra/core/host/container-apps-environment.bicep b/samples/scenarios/WorkItemFilteringSplitActivitiesPython/infra/core/host/container-apps-environment.bicep new file mode 100644 index 00000000..560d293c --- /dev/null +++ b/samples/scenarios/WorkItemFilteringSplitActivitiesPython/infra/core/host/container-apps-environment.bicep @@ -0,0 +1,27 @@ +metadata description = 'Creates an Azure Container Apps environment.' +param name string +param location string = resourceGroup().location +param tags object = {} + +@description('Subnet resource ID for the Container Apps environment') +param subnetResourceId string = '' + +@description('Whether to use an internal or external load balancer') +@allowed(['Internal', 'External']) +param loadBalancerType string = 'External' + +resource containerAppsEnvironment 'Microsoft.App/managedEnvironments@2023-05-01' = { + name: name + location: location + tags: tags + properties: { + vnetConfiguration: !empty(subnetResourceId) ? { + infrastructureSubnetId: subnetResourceId + internal: loadBalancerType == 'Internal' + } : null + } +} + +output defaultDomain string = containerAppsEnvironment.properties.defaultDomain +output id string = containerAppsEnvironment.id +output name string = containerAppsEnvironment.name diff --git a/samples/scenarios/WorkItemFilteringSplitActivitiesPython/infra/core/host/container-apps.bicep b/samples/scenarios/WorkItemFilteringSplitActivitiesPython/infra/core/host/container-apps.bicep new file mode 100644 index 00000000..45114237 --- /dev/null +++ b/samples/scenarios/WorkItemFilteringSplitActivitiesPython/infra/core/host/container-apps.bicep @@ -0,0 +1,44 @@ +metadata description = 'Creates an Azure Container Registry and an Azure Container Apps environment.' +param name string +param location string = resourceGroup().location +param tags object = {} + +param containerAppsEnvironmentName string +param containerRegistryName string +param containerRegistryAdminUserEnabled bool = false + +// Virtual network and subnet parameters +param subnetResourceId string = '' +param loadBalancerType string = 'External' + +module containerAppsEnvironment 'container-apps-environment.bicep' = { + name: '${name}-container-apps-environment' + params: { + name: containerAppsEnvironmentName + location: location + tags: tags + subnetResourceId: subnetResourceId + loadBalancerType: loadBalancerType + } +} + +module containerRegistry 'container-registry.bicep' = { + name: '${name}-container-registry' + params: { + name: containerRegistryName + location: location + adminUserEnabled: containerRegistryAdminUserEnabled + tags: tags + sku: { + name: 'Standard' + } + anonymousPullEnabled: false + } +} + +output defaultDomain string = containerAppsEnvironment.outputs.defaultDomain +output environmentName string = containerAppsEnvironment.outputs.name +output environmentId string = containerAppsEnvironment.outputs.id + +output registryLoginServer string = containerRegistry.outputs.loginServer +output registryName string = containerRegistry.outputs.name diff --git a/samples/scenarios/WorkItemFilteringSplitActivitiesPython/infra/core/host/container-registry.bicep b/samples/scenarios/WorkItemFilteringSplitActivitiesPython/infra/core/host/container-registry.bicep new file mode 100644 index 00000000..9ea04a2a --- /dev/null +++ b/samples/scenarios/WorkItemFilteringSplitActivitiesPython/infra/core/host/container-registry.bicep @@ -0,0 +1,134 @@ +metadata description = 'Creates an Azure Container Registry.' +param name string +param location string = resourceGroup().location +param tags object = {} + +@description('Indicates whether admin user is enabled') +param adminUserEnabled bool = false + +@description('Indicates whether anonymous pull is enabled') +param anonymousPullEnabled bool = false + +@description('Azure ad authentication as arm policy settings') +param azureADAuthenticationAsArmPolicy object = { + status: 'enabled' +} + +@description('Indicates whether data endpoint is enabled') +param dataEndpointEnabled bool = false + +@description('Encryption settings') +param encryption object = { + status: 'disabled' +} + +@description('Export policy settings') +param exportPolicy object = { + status: 'enabled' +} + +@description('Metadata search settings') +param metadataSearch string = 'Disabled' + +@description('Options for bypassing network rules') +param networkRuleBypassOptions string = 'AzureServices' + +@description('Public network access setting') +param publicNetworkAccess string = 'Enabled' + +@description('Quarantine policy settings') +param quarantinePolicy object = { + status: 'disabled' +} + +@description('Retention policy settings') +param retentionPolicy object = { + days: 7 + status: 'disabled' +} + +@description('Scope maps setting') +param scopeMaps array = [] + +@description('SKU settings') +param sku object = { + name: 'Basic' +} + +@description('Soft delete policy settings') +param softDeletePolicy object = { + retentionDays: 7 + status: 'disabled' +} + +@description('Trust policy settings') +param trustPolicy object = { + type: 'Notary' + status: 'disabled' +} + +@description('Zone redundancy setting') +param zoneRedundancy string = 'Disabled' + +@description('The log analytics workspace ID used for logging and monitoring') +param workspaceId string = '' + +// 2023-11-01-preview needed for metadataSearch +resource containerRegistry 'Microsoft.ContainerRegistry/registries@2023-11-01-preview' = { + name: name + location: location + tags: tags + sku: sku + properties: { + adminUserEnabled: adminUserEnabled + anonymousPullEnabled: anonymousPullEnabled + dataEndpointEnabled: dataEndpointEnabled + encryption: encryption + metadataSearch: metadataSearch + networkRuleBypassOptions: networkRuleBypassOptions + policies:{ + quarantinePolicy: quarantinePolicy + trustPolicy: trustPolicy + retentionPolicy: retentionPolicy + exportPolicy: exportPolicy + azureADAuthenticationAsArmPolicy: azureADAuthenticationAsArmPolicy + softDeletePolicy: softDeletePolicy + } + publicNetworkAccess: publicNetworkAccess + zoneRedundancy: zoneRedundancy + } + + resource scopeMap 'scopeMaps' = [for scopeMap in scopeMaps: { + name: scopeMap.name + properties: scopeMap.properties + }] +} + +resource diagnostics 'Microsoft.Insights/diagnosticSettings@2021-05-01-preview' = if (!empty(workspaceId)) { + name: 'registry-diagnostics' + scope: containerRegistry + properties: { + workspaceId: workspaceId + logs: [ + { + category: 'ContainerRegistryRepositoryEvents' + enabled: true + } + { + category: 'ContainerRegistryLoginEvents' + enabled: true + } + ] + metrics: [ + { + category: 'AllMetrics' + enabled: true + timeGrain: 'PT1M' + } + ] + } +} + +output id string = containerRegistry.id +output loginServer string = containerRegistry.properties.loginServer +output name string = containerRegistry.name diff --git a/samples/scenarios/WorkItemFilteringSplitActivitiesPython/infra/core/networking/vnet.bicep b/samples/scenarios/WorkItemFilteringSplitActivitiesPython/infra/core/networking/vnet.bicep new file mode 100644 index 00000000..57b5051d --- /dev/null +++ b/samples/scenarios/WorkItemFilteringSplitActivitiesPython/infra/core/networking/vnet.bicep @@ -0,0 +1,50 @@ +@description('The name of the Virtual Network') +param name string + +@description('The Azure region where the Virtual Network should exist') +param location string = resourceGroup().location + +@description('Optional tags for the resources') +param tags object = {} + +@description('The address prefixes of the Virtual Network') +param addressPrefixes array = ['10.0.0.0/16'] + +@description('The subnets to create in the Virtual Network') +param subnets array = [ + { + name: 'infrastructure-subnet' + properties: { + addressPrefix: '10.0.0.0/21' + delegations: [] + privateEndpointNetworkPolicies: 'Disabled' + privateLinkServiceNetworkPolicies: 'Enabled' + } + } + { + name: 'workload-subnet' + properties: { + addressPrefix: '10.0.8.0/21' + delegations: [] + privateEndpointNetworkPolicies: 'Disabled' + privateLinkServiceNetworkPolicies: 'Enabled' + } + } +] + +resource vnet 'Microsoft.Network/virtualNetworks@2022-07-01' = { + name: name + location: location + tags: tags + properties: { + addressSpace: { + addressPrefixes: addressPrefixes + } + subnets: subnets + } +} + +output id string = vnet.id +output name string = vnet.name +output infrastructureSubnetId string = resourceId('Microsoft.Network/virtualNetworks/subnets', name, 'infrastructure-subnet') +output workloadSubnetId string = resourceId('Microsoft.Network/virtualNetworks/subnets', name, 'workload-subnet') diff --git a/samples/scenarios/WorkItemFilteringSplitActivitiesPython/infra/core/security/registry-access.bicep b/samples/scenarios/WorkItemFilteringSplitActivitiesPython/infra/core/security/registry-access.bicep new file mode 100644 index 00000000..fc66837a --- /dev/null +++ b/samples/scenarios/WorkItemFilteringSplitActivitiesPython/infra/core/security/registry-access.bicep @@ -0,0 +1,19 @@ +metadata description = 'Assigns ACR Pull permissions to access an Azure Container Registry.' +param containerRegistryName string +param principalId string + +var acrPullRole = subscriptionResourceId('Microsoft.Authorization/roleDefinitions', '7f951dda-4ed3-4680-a7ca-43fe172d538d') + +resource aksAcrPull 'Microsoft.Authorization/roleAssignments@2022-04-01' = { + scope: containerRegistry // Use when specifying a scope that is different than the deployment scope + name: guid(subscription().id, resourceGroup().id, principalId, acrPullRole) + properties: { + roleDefinitionId: acrPullRole + principalType: 'ServicePrincipal' + principalId: principalId + } +} + +resource containerRegistry 'Microsoft.ContainerRegistry/registries@2023-01-01-preview' existing = { + name: containerRegistryName +} diff --git a/samples/scenarios/WorkItemFilteringSplitActivitiesPython/infra/core/security/role.bicep b/samples/scenarios/WorkItemFilteringSplitActivitiesPython/infra/core/security/role.bicep new file mode 100644 index 00000000..0b30cfd3 --- /dev/null +++ b/samples/scenarios/WorkItemFilteringSplitActivitiesPython/infra/core/security/role.bicep @@ -0,0 +1,21 @@ +metadata description = 'Creates a role assignment for a service principal.' +param principalId string + +@allowed([ + 'Device' + 'ForeignGroup' + 'Group' + 'ServicePrincipal' + 'User' +]) +param principalType string = 'ServicePrincipal' +param roleDefinitionId string + +resource role 'Microsoft.Authorization/roleAssignments@2022-04-01' = { + name: guid(subscription().id, resourceGroup().id, principalId, roleDefinitionId) + properties: { + principalId: principalId + principalType: principalType + roleDefinitionId: resourceId('Microsoft.Authorization/roleDefinitions', roleDefinitionId) + } +} diff --git a/samples/scenarios/WorkItemFilteringSplitActivitiesPython/infra/main.bicep b/samples/scenarios/WorkItemFilteringSplitActivitiesPython/infra/main.bicep new file mode 100644 index 00000000..5bcab161 --- /dev/null +++ b/samples/scenarios/WorkItemFilteringSplitActivitiesPython/infra/main.bicep @@ -0,0 +1,209 @@ +targetScope = 'subscription' + +@minLength(1) +@maxLength(64) +@description('Name of the the environment which is used to generate a short unique hash used in all resources.') +param environmentName string + +@minLength(1) +@description('Primary location for all resources') +param location string + +@description('Id of the user or app to assign application roles') +param principalId string = '' + +param containerAppsEnvName string = '' +param containerAppsAppName string = '' +param containerRegistryName string = '' +param dtsLocation string = 'centralus' +param dtsSkuName string = 'Consumption' +param dtsCapacity int = 1 +param dtsName string = '' +param taskHubName string = '' + +param clientServiceName string = 'client' +param orchestratorWorkerServiceName string = 'orchestrator-worker' +param validatorWorkerServiceName string = 'validator-worker' +param shipperWorkerServiceName string = 'shipper-worker' + +param resourceGroupName string = '' + +var abbrs = loadJsonContent('./abbreviations.json') + +var tags = { + 'azd-env-name': environmentName +} + +#disable-next-line no-unused-vars +var resourceToken = toLower(uniqueString(subscription().id, environmentName, location)) + +// Resource Group +resource rg 'Microsoft.Resources/resourceGroups@2021-04-01' = { + name: !empty(resourceGroupName) ? resourceGroupName : '${abbrs.resourcesResourceGroups}${environmentName}' + location: location + tags: tags +} + +// User-assigned managed identity for all container apps +module identity './app/user-assigned-identity.bicep' = { + name: 'identity' + scope: rg + params: { + name: 'dts-ca-identity' + } +} + +// Assign DTS Worker/Client role to the managed identity +module identityAssignDTS './core/security/role.bicep' = { + name: 'identityAssignDTS' + scope: rg + params: { + principalId: identity.outputs.principalId + roleDefinitionId: '0ad04412-c4d5-4796-b79c-f76d14c8d402' + principalType: 'ServicePrincipal' + } +} + +// Assign DTS role to the deploying user (for dashboard access) +module identityAssignDTSDash './core/security/role.bicep' = { + name: 'identityAssignDTSDash' + scope: rg + params: { + principalId: principalId + roleDefinitionId: '0ad04412-c4d5-4796-b79c-f76d14c8d402' + principalType: 'User' + } +} + +// Virtual network +module vnet './core/networking/vnet.bicep' = { + name: 'vnet' + scope: rg + params: { + name: '${abbrs.networkVirtualNetworks}${resourceToken}' + location: location + tags: tags + } +} + +// Container Apps Environment + Container Registry +module containerAppsEnv './core/host/container-apps.bicep' = { + name: 'container-apps' + scope: rg + params: { + name: 'app' + containerAppsEnvironmentName: !empty(containerAppsEnvName) ? containerAppsEnvName : '${abbrs.appManagedEnvironments}${resourceToken}' + containerRegistryName: !empty(containerRegistryName) ? containerRegistryName : '${abbrs.containerRegistryRegistries}${resourceToken}' + location: location + subnetResourceId: vnet.outputs.infrastructureSubnetId + loadBalancerType: 'External' + } +} + +// Durable Task Scheduler + Task Hub +module dts './app/dts.bicep' = { + scope: rg + name: 'dtsResource' + params: { + name: !empty(dtsName) ? dtsName : '${abbrs.dts}${resourceToken}' + taskhubname: !empty(taskHubName) ? taskHubName : '${abbrs.taskhub}${resourceToken}' + location: dtsLocation + tags: tags + ipAllowlist: ['0.0.0.0/0'] + skuName: dtsSkuName + skuCapacity: dtsCapacity + } +} + +// Client — schedules orchestrations and polls for results +module client 'app/app.bicep' = { + name: clientServiceName + scope: rg + params: { + appName: !empty(containerAppsAppName) ? '${containerAppsAppName}-client' : '${abbrs.appContainerApps}${resourceToken}-client' + containerAppsEnvironmentName: containerAppsEnv.outputs.environmentName + containerRegistryName: containerAppsEnv.outputs.registryName + userAssignedManagedIdentity: { + resourceId: identity.outputs.resourceId + clientId: identity.outputs.clientId + } + location: location + tags: tags + serviceName: 'client' + identityName: identity.outputs.name + dtsEndpoint: dts.outputs.dts_URL + taskHubName: dts.outputs.TASKHUB_NAME + } +} + +// Orchestrator Worker — handles orchestrations only +module orchestratorWorker 'app/app.bicep' = { + name: orchestratorWorkerServiceName + scope: rg + params: { + appName: !empty(containerAppsAppName) ? '${containerAppsAppName}-orchestrator' : '${abbrs.appContainerApps}${resourceToken}-orchestrator' + containerAppsEnvironmentName: containerAppsEnv.outputs.environmentName + containerRegistryName: containerAppsEnv.outputs.registryName + userAssignedManagedIdentity: { + resourceId: identity.outputs.resourceId + clientId: identity.outputs.clientId + } + location: location + tags: tags + serviceName: 'orchestrator-worker' + identityName: identity.outputs.name + dtsEndpoint: dts.outputs.dts_URL + taskHubName: dts.outputs.TASKHUB_NAME + workItemType: 'Orchestration' + } +} + +// Validator Worker — handles ValidateOrder activity only +module validatorWorker 'app/app.bicep' = { + name: validatorWorkerServiceName + scope: rg + params: { + appName: !empty(containerAppsAppName) ? '${containerAppsAppName}-validator' : '${abbrs.appContainerApps}${resourceToken}-validator' + containerAppsEnvironmentName: containerAppsEnv.outputs.environmentName + containerRegistryName: containerAppsEnv.outputs.registryName + userAssignedManagedIdentity: { + resourceId: identity.outputs.resourceId + clientId: identity.outputs.clientId + } + location: location + tags: tags + serviceName: 'validator-worker' + identityName: identity.outputs.name + dtsEndpoint: dts.outputs.dts_URL + taskHubName: dts.outputs.TASKHUB_NAME + workItemType: 'Activity' + } +} + +// Shipper Worker — handles ShipOrder activity only +module shipperWorker 'app/app.bicep' = { + name: shipperWorkerServiceName + scope: rg + params: { + appName: !empty(containerAppsAppName) ? '${containerAppsAppName}-shipper' : '${abbrs.appContainerApps}${resourceToken}-shipper' + containerAppsEnvironmentName: containerAppsEnv.outputs.environmentName + containerRegistryName: containerAppsEnv.outputs.registryName + userAssignedManagedIdentity: { + resourceId: identity.outputs.resourceId + clientId: identity.outputs.clientId + } + location: location + tags: tags + serviceName: 'shipper-worker' + identityName: identity.outputs.name + dtsEndpoint: dts.outputs.dts_URL + taskHubName: dts.outputs.TASKHUB_NAME + workItemType: 'Activity' + } +} + +output AZURE_LOCATION string = location +output AZURE_TENANT_ID string = tenant().tenantId +output AZURE_CONTAINER_REGISTRY_ENDPOINT string = containerAppsEnv.outputs.registryLoginServer +output AZURE_CONTAINER_REGISTRY_NAME string = containerAppsEnv.outputs.registryName +output AZURE_USER_ASSIGNED_IDENTITY_NAME string = identity.outputs.name diff --git a/samples/scenarios/WorkItemFilteringSplitActivitiesPython/infra/main.parameters.json b/samples/scenarios/WorkItemFilteringSplitActivitiesPython/infra/main.parameters.json new file mode 100644 index 00000000..c8d34538 --- /dev/null +++ b/samples/scenarios/WorkItemFilteringSplitActivitiesPython/infra/main.parameters.json @@ -0,0 +1,15 @@ +{ + "$schema": "https://schema.management.azure.com/schemas/2019-04-01/deploymentParameters.json#", + "contentVersion": "1.0.0.0", + "parameters": { + "environmentName": { + "value": "${AZURE_ENV_NAME}" + }, + "location": { + "value": "${AZURE_LOCATION}" + }, + "principalId": { + "value": "${AZURE_PRINCIPAL_ID}" + } + } +} diff --git a/samples/scenarios/WorkItemFilteringSplitActivitiesPython/run-local.sh b/samples/scenarios/WorkItemFilteringSplitActivitiesPython/run-local.sh new file mode 100644 index 00000000..08e2a23e --- /dev/null +++ b/samples/scenarios/WorkItemFilteringSplitActivitiesPython/run-local.sh @@ -0,0 +1,102 @@ +#!/usr/bin/env bash +# Runs the WorkItemFilteringSplitActivitiesPython sample locally. +# - Starts the DTS emulator (if not already running) +# - Creates a shared virtual environment and installs dependencies +# - Launches the three workers and the client, each in its own log file +# - Press Ctrl+C to stop everything + +set -euo pipefail + +SCRIPT_DIR="$(cd "$(dirname "${BASH_SOURCE[0]}")" && pwd)" +cd "$SCRIPT_DIR" + +LOG_DIR="$SCRIPT_DIR/.logs" +PID_FILE="$SCRIPT_DIR/.logs/pids" +VENV_DIR="$SCRIPT_DIR/.venv" +EMULATOR_NAME="dts-emulator" +EMULATOR_IMAGE="mcr.microsoft.com/dts/dts-emulator:latest" + +mkdir -p "$LOG_DIR" +: > "$PID_FILE" + +cleanup() { + echo "" + echo "[run-local] Stopping workers and client..." + if [[ -f "$PID_FILE" ]]; then + while read -r pid; do + if [[ -n "$pid" ]] && kill -0 "$pid" 2>/dev/null; then + kill "$pid" 2>/dev/null || true + fi + done < "$PID_FILE" + fi + if [[ "${KEEP_EMULATOR:-0}" != "1" ]]; then + echo "[run-local] Stopping DTS emulator container ($EMULATOR_NAME)..." + docker rm -f "$EMULATOR_NAME" >/dev/null 2>&1 || true + else + echo "[run-local] Leaving DTS emulator running (KEEP_EMULATOR=1)." + fi + echo "[run-local] Done." +} +trap cleanup EXIT INT TERM + +# 1. Ensure Docker is available +if ! command -v docker >/dev/null 2>&1; then + echo "[run-local] ERROR: Docker is required but not found in PATH." >&2 + exit 1 +fi + +# 2. Start the DTS emulator if needed +if docker ps --format '{{.Names}}' | grep -q "^${EMULATOR_NAME}$"; then + echo "[run-local] DTS emulator already running." +else + if docker ps -a --format '{{.Names}}' | grep -q "^${EMULATOR_NAME}$"; then + echo "[run-local] Removing stale emulator container..." + docker rm -f "$EMULATOR_NAME" >/dev/null + fi + echo "[run-local] Pulling DTS emulator image..." + docker pull "$EMULATOR_IMAGE" >/dev/null + echo "[run-local] Starting DTS emulator (dashboard: http://localhost:8082)..." + docker run -d --name "$EMULATOR_NAME" -p 8080:8080 -p 8082:8082 "$EMULATOR_IMAGE" >/dev/null +fi + +# 3. Create a shared virtual environment and install dependencies +if [[ ! -d "$VENV_DIR" ]]; then + echo "[run-local] Creating virtual environment at $VENV_DIR..." + python3 -m venv "$VENV_DIR" +fi +# shellcheck disable=SC1091 +source "$VENV_DIR/bin/activate" +echo "[run-local] Installing dependencies..." +pip install --quiet --upgrade pip +pip install --quiet -r src/client/requirements.txt + +# 4. Launch workers and client +start_proc() { + local name="$1" + local script="$2" + local log_file="$LOG_DIR/${name}.log" + echo "[run-local] Starting $name (logs: $log_file)" + python "$script" >"$log_file" 2>&1 & + echo $! >> "$PID_FILE" +} + +start_proc "orchestrator-worker" "src/orchestrator-worker/orchestrator_worker.py" +start_proc "validator-worker" "src/validator-worker/validator_worker.py" +start_proc "shipper-worker" "src/shipper-worker/shipper_worker.py" + +# Give workers a moment to connect before the client starts scheduling +sleep 3 + +start_proc "client" "src/client/client.py" + +echo "" +echo "[run-local] All processes started. Tailing logs (Ctrl+C to stop everything)..." +echo "[run-local] Logs are also saved under $LOG_DIR/" +echo "" + +# 5. Tail all logs until the user interrupts +tail -n +1 -F \ + "$LOG_DIR/orchestrator-worker.log" \ + "$LOG_DIR/validator-worker.log" \ + "$LOG_DIR/shipper-worker.log" \ + "$LOG_DIR/client.log" diff --git a/samples/scenarios/WorkItemFilteringSplitActivitiesPython/src/client/Dockerfile b/samples/scenarios/WorkItemFilteringSplitActivitiesPython/src/client/Dockerfile new file mode 100644 index 00000000..7ce6ba7b --- /dev/null +++ b/samples/scenarios/WorkItemFilteringSplitActivitiesPython/src/client/Dockerfile @@ -0,0 +1,9 @@ +FROM python:3.11-slim +WORKDIR /app + +COPY requirements.txt . +RUN pip install --no-cache-dir -r requirements.txt + +COPY . . + +CMD ["python", "client.py"] diff --git a/samples/scenarios/WorkItemFilteringSplitActivitiesPython/src/client/client.py b/samples/scenarios/WorkItemFilteringSplitActivitiesPython/src/client/client.py new file mode 100644 index 00000000..568fd249 --- /dev/null +++ b/samples/scenarios/WorkItemFilteringSplitActivitiesPython/src/client/client.py @@ -0,0 +1,140 @@ +import asyncio +import logging +import os +import time +from datetime import datetime, timedelta, timezone + +from azure.identity import DefaultAzureCredential +from durabletask import client as durable_client +from durabletask.azuremanaged.client import DurableTaskSchedulerClient + +logging.basicConfig( + level=logging.INFO, + format="%(asctime)s %(message)s", + datefmt="%H:%M:%S", +) +logger = logging.getLogger("Client") + +# Schedule a batch of orchestrations on a fixed interval so you can watch the +# workers scale over time. +ORCHESTRATIONS_PER_BATCH = 3 +INTERVAL_SECONDS = 30 +TOTAL_DURATION_MINUTES = 10 + + +async def main(): + logger.info("=== Work Item Filtering Demo — Client ===") + + taskhub_name = os.getenv("TASKHUB", "default") + endpoint = os.getenv("ENDPOINT", "http://localhost:8080") + managed_identity_client_id = os.getenv("AZURE_MANAGED_IDENTITY_CLIENT_ID") + + print(f"[Client] Using taskhub: {taskhub_name}") + print(f"[Client] Using endpoint: {endpoint}") + + # Use no credential for the local emulator, a user-assigned managed identity + # in Container Apps, or DefaultAzureCredential for local Azure development. + if endpoint == "http://localhost:8080": + credential = None + elif managed_identity_client_id: + credential = DefaultAzureCredential(managed_identity_client_id=managed_identity_client_id) + else: + credential = DefaultAzureCredential() + + client = DurableTaskSchedulerClient( + host_address=endpoint, + secure_channel=endpoint != "http://localhost:8080", + taskhub=taskhub_name, + token_credential=credential, + ) + + interval = timedelta(seconds=INTERVAL_SECONDS) + deadline = datetime.now(timezone.utc) + timedelta(minutes=TOTAL_DURATION_MINUTES) + + total_completed = 0 + total_failed = 0 + batch_number = 0 + + logger.info( + "Will schedule %d orchestrations every %ds for %d minutes.", + ORCHESTRATIONS_PER_BATCH, + INTERVAL_SECONDS, + TOTAL_DURATION_MINUTES, + ) + logger.info("(Make sure the Orchestrator, Validator, and Shipper workers are all running)\n") + + while datetime.now(timezone.utc) < deadline: + batch_number += 1 + logger.info("--- Batch #%d at %s ---", batch_number, datetime.now().strftime("%H:%M:%S")) + + instance_ids = [] + for i in range(1, ORCHESTRATIONS_PER_BATCH + 1): + order_id = f"ORD-B{batch_number:03d}-{i:03d}" + logger.info("Scheduling orchestration with orderId='%s'...", order_id) + instance_id = client.schedule_new_orchestration( + "order_processing_orchestrator", input=order_id + ) + instance_ids.append(instance_id) + logger.info(" -> Scheduled with InstanceId=%s", instance_id) + + # Wait for all orchestrations in this batch to complete. + batch_completed = 0 + batch_failed = 0 + for instance_id in instance_ids: + try: + state = client.wait_for_orchestration_completion(instance_id, timeout=120) + if state and state.runtime_status == durable_client.OrchestrationStatus.COMPLETED: + batch_completed += 1 + logger.info( + "COMPLETED | InstanceId=%s | Output: %s", + instance_id, + state.serialized_output, + ) + elif state: + batch_failed += 1 + logger.error( + "FAILED | InstanceId=%s | Status=%s | Error: %s", + instance_id, + state.runtime_status, + state.failure_details, + ) + except Exception as ex: # noqa: BLE001 + batch_failed += 1 + logger.error("Error waiting for orchestration %s: %s", instance_id, ex) + + total_completed += batch_completed + total_failed += batch_failed + logger.info( + "Batch #%d results: %d completed, %d failed", + batch_number, + batch_completed, + batch_failed, + ) + + # Wait for the next interval (unless we've passed the deadline). + now = datetime.now(timezone.utc) + if now < deadline: + remaining = deadline - now + wait_time = min(remaining, interval) + logger.info( + "Next batch in %.0fs (deadline in %.1f min)\n", + wait_time.total_seconds(), + remaining.total_seconds() / 60, + ) + await asyncio.sleep(wait_time.total_seconds()) + + logger.info( + "\n=== FINAL RESULTS: %d completed, %d failed across %d batches ===", + total_completed, + total_failed, + batch_number, + ) + + # Keep the process alive so Container Apps doesn't mark it as failed. + logger.info("Demo complete. Staying alive — press Ctrl+C to exit.") + while True: + time.sleep(3600) + + +if __name__ == "__main__": + asyncio.run(main()) diff --git a/samples/scenarios/WorkItemFilteringSplitActivitiesPython/src/client/requirements.txt b/samples/scenarios/WorkItemFilteringSplitActivitiesPython/src/client/requirements.txt new file mode 100644 index 00000000..1bc468dd --- /dev/null +++ b/samples/scenarios/WorkItemFilteringSplitActivitiesPython/src/client/requirements.txt @@ -0,0 +1,2 @@ +durabletask-azuremanaged +azure-identity diff --git a/samples/scenarios/WorkItemFilteringSplitActivitiesPython/src/orchestrator-worker/Dockerfile b/samples/scenarios/WorkItemFilteringSplitActivitiesPython/src/orchestrator-worker/Dockerfile new file mode 100644 index 00000000..d5e37b12 --- /dev/null +++ b/samples/scenarios/WorkItemFilteringSplitActivitiesPython/src/orchestrator-worker/Dockerfile @@ -0,0 +1,9 @@ +FROM python:3.11-slim +WORKDIR /app + +COPY requirements.txt . +RUN pip install --no-cache-dir -r requirements.txt + +COPY . . + +CMD ["python", "orchestrator_worker.py"] diff --git a/samples/scenarios/WorkItemFilteringSplitActivitiesPython/src/orchestrator-worker/orchestrator_worker.py b/samples/scenarios/WorkItemFilteringSplitActivitiesPython/src/orchestrator-worker/orchestrator_worker.py new file mode 100644 index 00000000..9440f091 --- /dev/null +++ b/samples/scenarios/WorkItemFilteringSplitActivitiesPython/src/orchestrator-worker/orchestrator_worker.py @@ -0,0 +1,112 @@ +import asyncio +import logging +import os + +from azure.identity import DefaultAzureCredential +from durabletask import task +from durabletask.azuremanaged.worker import DurableTaskSchedulerWorker + +logging.basicConfig( + level=logging.INFO, + format="%(asctime)s %(message)s", + datefmt="%H:%M:%S", +) +logger = logging.getLogger("Orchestrator") + +WORKER_NAME = "Orchestrator Worker" + + +def order_processing_orchestrator(ctx: task.OrchestrationContext, order_id: str): + """Calls two activities sequentially: + + 1. validate_order (routed to the Validator Worker) + 2. ship_order (routed to the Shipper Worker) + + Because each activity is registered in a different worker process, DTS routes + each activity work item to the correct worker via work item filtering. This + worker registers ONLY the orchestrator, so it never receives activity work items. + """ + if not ctx.is_replaying: + logger.info( + "[Orchestrator] Orchestration | Name=order_processing_orchestrator | " + "InstanceId=%s | Processing order '%s'", + ctx.instance_id, + order_id, + ) + + # Step 1: Validate the order (routed to the Validator Worker). + if not ctx.is_replaying: + logger.info( + "[Orchestrator] Orchestration | InstanceId=%s | Dispatching validate_order to Validator Worker...", + ctx.instance_id, + ) + validation_result = yield ctx.call_activity("validate_order", input=order_id) + + # Step 2: Ship the order (routed to the Shipper Worker). + if not ctx.is_replaying: + logger.info( + "[Orchestrator] Orchestration | InstanceId=%s | Dispatching ship_order to Shipper Worker...", + ctx.instance_id, + ) + shipping_result = yield ctx.call_activity("ship_order", input=order_id) + + combined = ( + f"Order '{order_id}' => Validation: [{validation_result}], " + f"Shipping: [{shipping_result}]" + ) + + if not ctx.is_replaying: + logger.info( + "[Orchestrator] Orchestration | InstanceId=%s | Completed: %s", + ctx.instance_id, + combined, + ) + + return combined + + +async def main(): + taskhub_name = os.getenv("TASKHUB", "default") + endpoint = os.getenv("ENDPOINT", "http://localhost:8080") + managed_identity_client_id = os.getenv("AZURE_MANAGED_IDENTITY_CLIENT_ID") + + print(f"[{WORKER_NAME}] Using taskhub: {taskhub_name}") + print(f"[{WORKER_NAME}] Using endpoint: {endpoint}") + + # Use no credential for the local emulator, a user-assigned managed identity + # in Container Apps, or DefaultAzureCredential for local Azure development. + if endpoint == "http://localhost:8080": + credential = None + elif managed_identity_client_id: + credential = DefaultAzureCredential(managed_identity_client_id=managed_identity_client_id) + else: + credential = DefaultAzureCredential() + + with DurableTaskSchedulerWorker( + host_address=endpoint, + secure_channel=endpoint != "http://localhost:8080", + taskhub=taskhub_name, + token_credential=credential, + ) as worker: + + # Register ONLY the orchestrator — no activities. + worker.add_orchestrator(order_processing_orchestrator) + + # Auto-generate work item filters from the registry, so this worker + # receives ONLY orchestration work items — never activity work items. + worker.use_work_item_filters() + + worker.start() + logger.info("[%s] Ready — processing orchestrations only.", WORKER_NAME) + + try: + while True: + await asyncio.sleep(1) + except KeyboardInterrupt: + logger.info("[%s] Shutdown initiated", WORKER_NAME) + + logger.info("[%s] Stopped", WORKER_NAME) + + +if __name__ == "__main__": + asyncio.run(main()) diff --git a/samples/scenarios/WorkItemFilteringSplitActivitiesPython/src/orchestrator-worker/requirements.txt b/samples/scenarios/WorkItemFilteringSplitActivitiesPython/src/orchestrator-worker/requirements.txt new file mode 100644 index 00000000..1bc468dd --- /dev/null +++ b/samples/scenarios/WorkItemFilteringSplitActivitiesPython/src/orchestrator-worker/requirements.txt @@ -0,0 +1,2 @@ +durabletask-azuremanaged +azure-identity diff --git a/samples/scenarios/WorkItemFilteringSplitActivitiesPython/src/shipper-worker/Dockerfile b/samples/scenarios/WorkItemFilteringSplitActivitiesPython/src/shipper-worker/Dockerfile new file mode 100644 index 00000000..d748b311 --- /dev/null +++ b/samples/scenarios/WorkItemFilteringSplitActivitiesPython/src/shipper-worker/Dockerfile @@ -0,0 +1,9 @@ +FROM python:3.11-slim +WORKDIR /app + +COPY requirements.txt . +RUN pip install --no-cache-dir -r requirements.txt + +COPY . . + +CMD ["python", "shipper_worker.py"] diff --git a/samples/scenarios/WorkItemFilteringSplitActivitiesPython/src/shipper-worker/requirements.txt b/samples/scenarios/WorkItemFilteringSplitActivitiesPython/src/shipper-worker/requirements.txt new file mode 100644 index 00000000..1bc468dd --- /dev/null +++ b/samples/scenarios/WorkItemFilteringSplitActivitiesPython/src/shipper-worker/requirements.txt @@ -0,0 +1,2 @@ +durabletask-azuremanaged +azure-identity diff --git a/samples/scenarios/WorkItemFilteringSplitActivitiesPython/src/shipper-worker/shipper_worker.py b/samples/scenarios/WorkItemFilteringSplitActivitiesPython/src/shipper-worker/shipper_worker.py new file mode 100644 index 00000000..97e6bef8 --- /dev/null +++ b/samples/scenarios/WorkItemFilteringSplitActivitiesPython/src/shipper-worker/shipper_worker.py @@ -0,0 +1,85 @@ +import asyncio +import logging +import os +import random + +from azure.identity import DefaultAzureCredential +from durabletask import task +from durabletask.azuremanaged.worker import DurableTaskSchedulerWorker + +logging.basicConfig( + level=logging.INFO, + format="%(asctime)s %(message)s", + datefmt="%H:%M:%S", +) +logger = logging.getLogger("Shipper") + +WORKER_NAME = "Shipper Worker" + + +def ship_order(ctx: task.ActivityContext, order_id: str) -> str: + """Ships an order. Registered ONLY in the Shipper Worker, + so DTS routes ship_order work items exclusively to this worker. + """ + logger.info( + "[Shipper] Activity | Name=ship_order | InstanceId=%s | Shipping order '%s'...", + ctx.orchestration_id, + order_id, + ) + + tracking_number = f"TRACK-{order_id}-{random.randint(1000, 9999)}" + result = f"Shipped with tracking {tracking_number}" + + logger.info( + "[Shipper] Activity | Name=ship_order | InstanceId=%s | Result: %s", + ctx.orchestration_id, + result, + ) + return result + + +async def main(): + taskhub_name = os.getenv("TASKHUB", "default") + endpoint = os.getenv("ENDPOINT", "http://localhost:8080") + managed_identity_client_id = os.getenv("AZURE_MANAGED_IDENTITY_CLIENT_ID") + + print(f"[{WORKER_NAME}] Using taskhub: {taskhub_name}") + print(f"[{WORKER_NAME}] Using endpoint: {endpoint}") + + # Use no credential for the local emulator, a user-assigned managed identity + # in Container Apps, or DefaultAzureCredential for local Azure development. + if endpoint == "http://localhost:8080": + credential = None + elif managed_identity_client_id: + credential = DefaultAzureCredential(managed_identity_client_id=managed_identity_client_id) + else: + credential = DefaultAzureCredential() + + with DurableTaskSchedulerWorker( + host_address=endpoint, + secure_channel=endpoint != "http://localhost:8080", + taskhub=taskhub_name, + token_credential=credential, + ) as worker: + + # Register ONLY the ship_order activity — no orchestrations. + worker.add_activity(ship_order) + + # Auto-generate work item filters from the registry, so this worker + # receives ONLY ship_order activity work items. + worker.use_work_item_filters() + + worker.start() + logger.info("[%s] Ready — processing ship_order activity only.", WORKER_NAME) + + try: + while True: + await asyncio.sleep(1) + except KeyboardInterrupt: + logger.info("[%s] Shutdown initiated", WORKER_NAME) + + logger.info("[%s] Stopped", WORKER_NAME) + + +if __name__ == "__main__": + asyncio.run(main()) diff --git a/samples/scenarios/WorkItemFilteringSplitActivitiesPython/src/validator-worker/Dockerfile b/samples/scenarios/WorkItemFilteringSplitActivitiesPython/src/validator-worker/Dockerfile new file mode 100644 index 00000000..8c77feba --- /dev/null +++ b/samples/scenarios/WorkItemFilteringSplitActivitiesPython/src/validator-worker/Dockerfile @@ -0,0 +1,9 @@ +FROM python:3.11-slim +WORKDIR /app + +COPY requirements.txt . +RUN pip install --no-cache-dir -r requirements.txt + +COPY . . + +CMD ["python", "validator_worker.py"] diff --git a/samples/scenarios/WorkItemFilteringSplitActivitiesPython/src/validator-worker/requirements.txt b/samples/scenarios/WorkItemFilteringSplitActivitiesPython/src/validator-worker/requirements.txt new file mode 100644 index 00000000..1bc468dd --- /dev/null +++ b/samples/scenarios/WorkItemFilteringSplitActivitiesPython/src/validator-worker/requirements.txt @@ -0,0 +1,2 @@ +durabletask-azuremanaged +azure-identity diff --git a/samples/scenarios/WorkItemFilteringSplitActivitiesPython/src/validator-worker/validator_worker.py b/samples/scenarios/WorkItemFilteringSplitActivitiesPython/src/validator-worker/validator_worker.py new file mode 100644 index 00000000..ef707492 --- /dev/null +++ b/samples/scenarios/WorkItemFilteringSplitActivitiesPython/src/validator-worker/validator_worker.py @@ -0,0 +1,83 @@ +import asyncio +import logging +import os + +from azure.identity import DefaultAzureCredential +from durabletask import task +from durabletask.azuremanaged.worker import DurableTaskSchedulerWorker + +logging.basicConfig( + level=logging.INFO, + format="%(asctime)s %(message)s", + datefmt="%H:%M:%S", +) +logger = logging.getLogger("Validator") + +WORKER_NAME = "Validator Worker" + + +def validate_order(ctx: task.ActivityContext, order_id: str) -> str: + """Validates an incoming order. Registered ONLY in the Validator Worker, + so DTS routes validate_order work items exclusively to this worker. + """ + logger.info( + "[Validator] Activity | Name=validate_order | InstanceId=%s | Validating order '%s'...", + ctx.orchestration_id, + order_id, + ) + + result = f"Order {order_id} is valid" + + logger.info( + "[Validator] Activity | Name=validate_order | InstanceId=%s | Result: %s", + ctx.orchestration_id, + result, + ) + return result + + +async def main(): + taskhub_name = os.getenv("TASKHUB", "default") + endpoint = os.getenv("ENDPOINT", "http://localhost:8080") + managed_identity_client_id = os.getenv("AZURE_MANAGED_IDENTITY_CLIENT_ID") + + print(f"[{WORKER_NAME}] Using taskhub: {taskhub_name}") + print(f"[{WORKER_NAME}] Using endpoint: {endpoint}") + + # Use no credential for the local emulator, a user-assigned managed identity + # in Container Apps, or DefaultAzureCredential for local Azure development. + if endpoint == "http://localhost:8080": + credential = None + elif managed_identity_client_id: + credential = DefaultAzureCredential(managed_identity_client_id=managed_identity_client_id) + else: + credential = DefaultAzureCredential() + + with DurableTaskSchedulerWorker( + host_address=endpoint, + secure_channel=endpoint != "http://localhost:8080", + taskhub=taskhub_name, + token_credential=credential, + ) as worker: + + # Register ONLY the validate_order activity — no orchestrations. + worker.add_activity(validate_order) + + # Auto-generate work item filters from the registry, so this worker + # receives ONLY validate_order activity work items. + worker.use_work_item_filters() + + worker.start() + logger.info("[%s] Ready — processing validate_order activity only.", WORKER_NAME) + + try: + while True: + await asyncio.sleep(1) + except KeyboardInterrupt: + logger.info("[%s] Shutdown initiated", WORKER_NAME) + + logger.info("[%s] Stopped", WORKER_NAME) + + +if __name__ == "__main__": + asyncio.run(main())