using Dapper; //using Hangfire; using MisIngredientesVue.Core.Models; using System; using System.Collections.Generic; using System.IO; using System.Linq; using System.Net; using System.Net.Security; using System.Text; using System.Threading.Tasks; namespace MisIngredientesVue.Core.Services { public interface IWebHookService { Task> GetWebhooks(string tenantId); Task AddWebHooks(string tenantId, string url, params string[] events); Task> GetWebhookUrls(string tenantId, string eventName); Task TriggerWebHooks(string tenantId, string eventName, string payload, string[] urls = null); } public class ResponseWebHook { public string Name { get; set; } public string[] Events { get; set; } } public class WebHookService : IWebHookService { private readonly IRepository _repository; private readonly ILogRepository _logRepository; public WebHookService(IRepository repository, ILogRepository logRepository) { _repository = repository; _logRepository = logRepository; } public async Task AddWebHooks(string tenantId, string url, params string[] events) { await _repository.Query(async (db) => { var queryBuilder = new List(); var parameters = new DynamicParameters(); parameters.Add("tenantId", tenantId); parameters.Add("url", url); var index = 0; foreach(var ev in events) { var query = @$"INSERT INTO public.webhooks (tenantid, url, eventname) VALUES (@tenantId, @url, @ev{index}) ON CONFLICT(tenantid, url, eventname) DO UPDATE SET eventname = @ev{index};"; queryBuilder.Add(query); parameters.Add($"ev{index}", ev); index++; } await db.QueryAsync(string.Join(" ", queryBuilder), parameters); return null; }); } public async Task> GetWebhooks(string tenantId) { return await _repository.Query(async (db) => { var query = @$"SELECT * FROM public.webhooks WHERE tenantId = @tenantId"; var webhooks = await db.QueryAsync(query, new { tenantId }); return webhooks.GroupBy(x => x.Url).Select(x => new ResponseWebHook() { Name = x.Key, Events = x.Select(y => y.EventName).ToArray() }); }); } public async Task TriggerWebHooks(string tenantId, string eventName, string payload, string[] urls = null) { try { //Console.WriteLine($"[TriggerWebHooks] {eventName} = {payload}"); await TriggerWebHooksAsync(tenantId, eventName, payload, urls); //BackgroundJob.Enqueue(() => TriggerWebHooksAsync(tenantId, eventName, payload)); } catch (Exception ex) { Console.WriteLine($"[ERROR TriggerWebHooks] " + ex.Message); } } public async Task> GetWebhookUrls(string tenantId, string eventName) { var urls = await _repository.Query>(async (db) => { var query = @$"SELECT DISTINCT url FROM public.webhooks WHERE tenantId = @tenantId AND eventname = @eventName"; return await db.QueryAsync(query, new { tenantId, eventName }); }); return urls.Where(x => !string.IsNullOrEmpty(x)); } public async Task TriggerWebHooksAsync(string tenantId, string eventName, string payload, string[] urls = null) { if (urls == null) urls = (await GetWebhookUrls(tenantId, eventName)).ToArray(); foreach (var url in urls) { string eventId = Guid.NewGuid().ToString(); try { if (string.IsNullOrEmpty(url)) continue; var body = $"{{ EventId: \"{eventId}\", EventType: \"{eventName}\", Payload: {payload} }}"; await HttpRequest(url, body, "POST", ""); //await _logRepository.InsertAsync(new WebHookCall() { // DateTime = DateTime.UtcNow, // EventName = eventName, // Request = payload, // Response = "", // Success = true, // TenantId = tenantId, // Url = url, // EventId = eventId //}); } catch (Exception ex) { //await _logRepository.InsertAsync(new WebHookCall() //{ // DateTime = DateTime.UtcNow, // EventName = eventName, // Request = payload, // Response = ex.Message, // Success = false, // TenantId = tenantId, // Url = url, // EventId = eventId //}); } } } private async Task HttpRequest(string url, string body, string method, string authHeader) { ServicePointManager.ServerCertificateValidationCallback += (sender, certificate, chain, errors) => { // local dev, just approve all certs return true; return errors == SslPolicyErrors.None; }; var httpWebRequest = (HttpWebRequest)WebRequest.Create(url); httpWebRequest.ContentType = "application/json"; httpWebRequest.Method = method; var response = string.Empty; if (!string.IsNullOrEmpty(body)) { using (var req = httpWebRequest.GetRequestStream()) { var bodyBytes = Encoding.UTF8.GetBytes(body); await req.WriteAsync(bodyBytes, 0, bodyBytes.Length); } } try { using (var res = await httpWebRequest.GetResponseAsync()) { using (var str = res.GetResponseStream()) { } } } catch (WebException ex) { if (ex.Response == null) throw new Exception(ex.Message); using (var res = ex.Response) { using (var str = res.GetResponseStream()) { using (var reader = new StreamReader(str)) { response = await reader.ReadToEndAsync(); if(!string.IsNullOrWhiteSpace(response)) throw new Exception(response); throw new Exception(ex.Message); } } } } } } }