-
Notifications
You must be signed in to change notification settings - Fork 0
Expand file tree
/
Copy pathFuncQueueListener.cs
More file actions
112 lines (98 loc) · 4.99 KB
/
Copy pathFuncQueueListener.cs
File metadata and controls
112 lines (98 loc) · 4.99 KB
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
using System;
using System.Text;
using System.Text.Json;
using Azure.Storage.Queues;
using Azure.Storage.Queues.Models;
using FuncTriggerManagerSvc.Models;
using FuncTriggerManagerSvc.Utils;
using Microsoft.Azure.Functions.Worker;
using Microsoft.Extensions.Configuration;
using Microsoft.Extensions.Logging;
using Newtonsoft.Json;
namespace FuncTriggerManagerSvc
{
public class FuncQueueListener
{
private readonly ILogger<FuncQueueListener> _logger;
private readonly IConfiguration _configuration;
public FuncQueueListener(ILogger<FuncQueueListener> logger, IConfiguration configuration)
{
_logger = logger;
_configuration = configuration;
}
[Function(nameof(FuncQueueListener))]
public async Task Run(
[QueueTrigger("%QueueName%", Connection = "StorageConnection")] QueueMessage queueMessage)
{
_logger.LogInformation($"C# Queue trigger function processed: {queueMessage.MessageText}");
// Get the connection string from configuration
string? connectionString = _configuration.GetValue<string>("StorageConnection");
if (string.IsNullOrEmpty(connectionString))
{
_logger.LogError("StorageConnection is not configured or is null/empty.");
throw new InvalidOperationException("StorageConnection configuration is missing.");
}
string? subscriptionId = _configuration.GetValue<string>("AzureSubscriptionId");
if (string.IsNullOrEmpty(subscriptionId))
{
_logger.LogError("Azure Subscription Id is not configured or is null/empty.");
throw new InvalidOperationException("AzureSubscriptionId configuration is missing.");
}
string? queueName = _configuration.GetValue<string>("QueueName");
if (string.IsNullOrEmpty(queueName))
{
_logger.LogError("QueueName is not configured or is null/empty.");
throw new InvalidOperationException("QueueName configuration is missing.");
}
// Access the message conten
string messageContent = queueMessage.MessageText;
var shortCircuitMsg = new FuncTriggerMsg();
_logger.LogInformation($"Message content: {messageContent}");
try
{
shortCircuitMsg = JsonConvert.DeserializeObject<FuncTriggerMsg>(messageContent);
_logger.LogInformation($"Deserialized message for function: {shortCircuitMsg!.FunctionAppName}/{shortCircuitMsg!.FunctionName}");
}
catch (Newtonsoft.Json.JsonException ex)
{
_logger.LogError(ex, "Failed to deserialize message content.");
throw;
}
// Process the message to either disable or enable the function app.
try
{
// Get the function app details
var functionApp = await WebAppConfigurator.GetFunctionAsync(subscriptionId!, shortCircuitMsg.ResourceGroupName,
shortCircuitMsg.FunctionAppName, _logger);
// set the function status based on the message
await WebAppConfigurator.FunctionConditionAsync(functionApp, shortCircuitMsg.FunctionName,
shortCircuitMsg.DisableFunction, _logger);
// If disable function is true, republish the message to the queue if the disable period is greater than 0
if (shortCircuitMsg.DisableFunction && shortCircuitMsg.DisablePeriodMinutes > 0)
{
_logger.LogInformation($"Function {shortCircuitMsg.FunctionName} is disabled for {shortCircuitMsg.DisablePeriodMinutes} minutes");
// Set the visibility timeout to the specified period
TimeSpan visibilityTimeout = TimeSpan.FromMinutes(shortCircuitMsg.DisablePeriodMinutes);
shortCircuitMsg.DisableFunction = false;
// Create a queue client
QueueClient queueClient = new QueueClient(connectionString, queueName);
await queueClient.CreateIfNotExistsAsync();
// Send the message with a custom visibility timeout
var messageOut = JsonConvert.SerializeObject(shortCircuitMsg);
string base64Message = Convert.ToBase64String(Encoding.UTF8.GetBytes(messageOut));
await queueClient.SendMessageAsync(base64Message, visibilityTimeout);
_logger.LogInformation($"Message republished to output queue with visibility timeout of {visibilityTimeout}");
}
else
{
_logger.LogInformation($"Function {shortCircuitMsg.FunctionName} is enabled");
}
}
catch (Exception ex)
{
_logger.LogError(ex, "Failed to process the message");
throw;
}
}
}
}