anyxchange.streaming.prepareListingModerationPipeline
pipeline klog tables adminStash adminEmails moderationServiceOpt =
do
use DlqFailureStage DlqUnknownStage
use Failure message
listingDlqTbl = DeploymentTables.listingModerationDlqTable tables
processEvent lid event =
(ListingId lidText) = lid
smokeListingPrefix = "moderation-listing-smoke-token-"
if startsWith smokeListingPrefix lidText then
smokeTokenText = Text.drop (Text.size smokeListingPrefix) lidText
toRemote do
_ =
recordPipelineProbeConsumption
tables
moderationListingPipelineSmokeConsumedTokenKey
smokeTokenText
+ ()
info
"moderation.listing-pipeline.smoke-klog-consumed"
[("listingId", lidText), ("token", smokeTokenText)]
-- Signature
anyxchange.streaming.prepareListingModerationPipeline :
Pipeline
-> KLog ListingId ListingModerationEvent
-> DeploymentTables
-> Stash [schema_5_2_0.UserId]
-> [Text]
-> Optional (ServiceHash ModerationRequest ModerationResult)
-> '{Exception, Cloud} ()
if startsWith smokeListingPrefix lidTxt then
smokeTokenText = Text.drop (Text.size smokeListingPrefix) lidTxt
toRemote do
_ =
recordPipelineProbeConsumption
- tables matchingPipelineSmokeConsumedTokenKey smokeTokenText
+ tables matchingPipelineSmokeConsumedTokenKey smokeTokenText ()
info
"pipeline.smoke-klog-consumed"
[("listingId", lidTxt), ("token", smokeTokenText)]
else
-- Signature
anyxchange.monitoring.recordPipelineProbeConsumption :
DeploymentTables -> Text -> Text -> '{Remote} ()
if startsWith smokeUserPrefix uidText then
smokeTokenText = Text.drop (Text.size smokeUserPrefix) uidText
toRemote do
_ =
recordPipelineProbeConsumption
- tables moderationUserPipelineSmokeConsumedTokenKey smokeTokenText
+ tables
+ moderationUserPipelineSmokeConsumedTokenKey
+ smokeTokenText
+ ()
info
"moderation.user-pipeline.smoke-klog-consumed"
[("userId", uidText), ("token", smokeTokenText)]
else