using System; using System.Collections.Generic; using System.Linq; using System.Threading.Tasks; using MySqlConnector; using Sony.Filtr.Database; using StackExchange.Redis; using StackExchange.Redis.Extensions.Core.Implementations; using Sony.Filtr.Utility.Extensions; namespace Sony.Filtr.Tasks { public class ScheduledTaskManager { private IDatabaseAsync _redis; public ScheduledTaskManager(IDatabaseAsync redis) { _redis = redis; } public async Task LogTaskStartedAsync(ScheduledTask scheduledTask, Guid scheduledTaskLogId) { using (var conn = await DatabaseHandler.GetOpenConnectionAsync()) { var cmd = new MySqlCommand("INSERT INTO tblScheduledTaskLog (Id, ScheduledTaskId, Started) VALUES (@Id, @scheduledTaskId, @started)", conn); cmd.Parameters.AddWithValue("@Id", scheduledTaskLogId); cmd.Parameters.AddWithValue("@scheduledTaskId", scheduledTask.Id); cmd.Parameters.AddWithValue("@started", DateTime.UtcNow); await cmd.ExecuteNonQueryAsync(); } } public async Task LogTaskFinishedAsync(ScheduledTask scheduledTask) { await _redis.LockReleaseAsync("Scheduled_Task_Locked_" + scheduledTask.Name, scheduledTask.Id); } public async Task IsTaskCurrentlyRunById(ScheduledTask scheduledTask) { return !await _redis.LockTakeAsync("Scheduled_Task_Locked_" + scheduledTask.Name, scheduledTask.Id, TimeSpan.FromHours(24), CommandFlags.None); } public async Task LogTaskFinishedAsync(Guid scheduledTaskLogId, ScheduledTaskLog log) { using (var conn = await DatabaseHandler.GetOpenConnectionAsync()) { var cmd = new MySqlCommand("UPDATE tblScheduledTaskLog SET finished=@finished, Error=@error, ErrorMessage=@errorMessage WHERE Id=@scheduledTaskLogId", conn); cmd.Parameters.AddWithValue("@scheduledTaskLogId", scheduledTaskLogId); cmd.Parameters.AddWithValue("@finished", DateTime.UtcNow); cmd.Parameters.AddWithValue("@error", log.Error); cmd.Parameters.AddWithValue("@errorMessage", log.ErrorMessage); await cmd.ExecuteNonQueryAsync(); } } public async Task AddTaskDataLogAsync(ScheduledTaskDataLog dataLog) { using (var conn = await DatabaseHandler.GetOpenConnectionAsync()) { var dataLogCmd = new MySqlCommand("INSERT INTO tblScheduledTaskDataLog (ScheduledTaskLogId, DataId, Started, Finished, RepresentationDate, DataSource, Rows, FailedRows, Error, ErrorMessage) " + "VALUES (@ScheduledTaskLogId, @dataId, @started, @finished, @representationDate, @dataSource, @rows, @failedRows, @error, @errorMessage) ", conn); dataLogCmd.Parameters.AddWithValue("@scheduledTaskLogId", dataLog.ScheduledTaskLogId); dataLogCmd.Parameters.AddWithValue("@dataId", dataLog.DataId); dataLogCmd.Parameters.AddWithValue("@started", dataLog.Started); dataLogCmd.Parameters.AddWithValue("@Finished", dataLog.Finished); dataLogCmd.Parameters.AddWithValue("@representationDate", dataLog.RepresentationDate); dataLogCmd.Parameters.AddWithValue("@dataSource", dataLog.DataSource); dataLogCmd.Parameters.AddWithValue("@rows", dataLog.Rows); dataLogCmd.Parameters.AddWithValue("@failedRows", dataLog.FailedRows); dataLogCmd.Parameters.AddWithValue("@error", dataLog.Error); dataLogCmd.Parameters.AddWithValue("@errorMessage", dataLog.ErrorMessage); await dataLogCmd.ExecuteNonQueryAsync(); dataLog.Id = (int)dataLogCmd.LastInsertedId; } return dataLog; } public async Task UpdateTaskDataLogAsync(ScheduledTaskDataLog dataLog) { using (var conn = await DatabaseHandler.GetOpenConnectionAsync()) { var dataLogCmd = new MySqlCommand("UPDATE tblScheduledTaskDataLog SET Started = @started, Finished = @finished, RepresentationDate = @representationDate, DataSource = @dataSource, " + "Rows = @rows, FailedRows = @failedRows, Error = @error, ErrorMessage = @errorMessage " + "WHERE Id = @dataLogId", conn); dataLogCmd.Parameters.AddWithValue("@dataLogId", dataLog.Id); dataLogCmd.Parameters.AddWithValue("@started", dataLog.Started); dataLogCmd.Parameters.AddWithValue("@Finished", dataLog.Finished); dataLogCmd.Parameters.AddWithValue("@representationDate", dataLog.RepresentationDate); dataLogCmd.Parameters.AddWithValue("@dataSource", dataLog.DataSource); dataLogCmd.Parameters.AddWithValue("@rows", dataLog.Rows); dataLogCmd.Parameters.AddWithValue("@failedRows", dataLog.FailedRows); dataLogCmd.Parameters.AddWithValue("@error", dataLog.Error); dataLogCmd.Parameters.AddWithValue("@errorMessage", dataLog.ErrorMessage); await dataLogCmd.ExecuteNonQueryAsync(); } } public async Task GetOrAddScheduledTaskAsync(string key) { ScheduledTask task; using (var conn = await DatabaseHandler.GetOpenConnectionAsync()) { var cmd = new MySqlCommand("SELECT Id, Name FROM tblScheduledTask WHERE ParameterId=@parameterId", conn); cmd.Parameters.AddWithValue("@parameterId", key); var reader = await cmd.ExecuteReaderAsync(); if (await reader.ReadAsync()) { var id = reader.GetInt32(0); var name = reader.GetString(1); task = new ScheduledTask { Id = id, Name = name }; reader.Close(); return task; } reader.Close(); var addCmd = new MySqlCommand("INSERT INTO tblScheduledTask (Name, ParameterId) VALUES (@parameterId, @parameterId)", conn); addCmd.Parameters.AddWithValue("@parameterId", key); await addCmd.ExecuteNonQueryAsync(); task = new ScheduledTask() { Id = (int)addCmd.LastInsertedId, Name = key }; } return task; } public async Task GetLatestFinishedLogAsync(string scheduledTaskName) { using (var conn = await DatabaseHandler.GetOpenReadOnlyConnectionAsync()) { using(var cmd = new MySqlCommand("SELECT tl.Finished, tl.Error, tl.ErrorMessage FROM tblScheduledTaskLog AS tl " + "INNER JOIN tblScheduledTask AS t ON tl.ScheduledTaskId = t.id " + "WHERE t.name = @scheduledTaskName AND tl.finished IS NOT NULL AND Error = 0 " + "ORDER BY Finished DESC " + "LIMIT 1", conn)) { cmd.Parameters.AddWithValue("@scheduledTaskName", scheduledTaskName); using(var reader = await cmd.ExecuteReaderAsync()) { if (reader.Read()) { var finished = reader.GetUtcDateTime("Finished"); var error = reader.GetBoolean("Error"); var errorMessage = reader.GetSafeString("ErrorMessage"); return new ScheduledTaskLog() { Finished = finished, Error = error, ErrorMessage = errorMessage }; } } } } return null; } } public class ScheduledTask { public int Id { get; set; } public string Name { get; set; } } public class ScheduledTaskLog { public ScheduledTaskLog() { DataLogs = new List(); } public DateTimeOffset? Finished { get; set; } public bool Error { get; set; } public string ErrorMessage { get; set; } public List DataLogs { get; set; } } public class ScheduledTaskDataLog { public ScheduledTaskDataLog(string dataId, DateTimeOffset started, Guid scheduledTaskLogId) { DataId = dataId; Started = started; ScheduledTaskLogId = scheduledTaskLogId; } public int Id { get; set; } public Guid ScheduledTaskLogId { get; set; } public string DataId { get; set; } public DateTimeOffset Started { get; set; } public DateTimeOffset Finished { get; set; } public DateTime RepresentationDate { get; set; } public string DataSource { get; set; } public int Rows { get; set; } public int FailedRows { get; set; } public bool Error { get; set; } public string ErrorMessage { get; set; } } public class ScheduledTasksInfo { public ScheduledTasksInfo() { } public ScheduledTasksInfo(HashEntry[] hash) { var hashDic = hash.ToDictionary(k => k.Name, v => v.Value); var properties = GetType().GetProperties().Where(p=> p.CanWrite && IsValidType(p.PropertyType)); foreach (var property in properties) { if (!hashDic.ContainsKey(property.Name)) continue; var value = hashDic[property.Name]; if (property.PropertyType == typeof (DateTime?)) { DateTime dateTime; if (DateTime.TryParse(value, out dateTime)) { property.SetValue(this, dateTime); } } else if (property.PropertyType == typeof (int?)) { int integer; if (int.TryParse(value, out integer)) { property.SetValue(this, integer); } } } } private bool IsValidType(Type type) { return (type == typeof (DateTime?) || type == typeof (int?)); } public HashEntry[] GetHashFields() { var properties = GetType().GetProperties().Where(p => p.CanRead && IsValidType(p.PropertyType) && p.GetValue(this) != null); return properties.Select(p => new HashEntry(p.Name, p.GetValue(this).ToString())).ToArray(); } public int? UniquePlaylists { get; set; } public int? UniqueEditorialPlaylists { get; set; } public int? UniqueBuzzPlaylists { get; set; } public DateTime? UpdatePlaylistsStart { get; set; } public DateTime? UpdatePlaylistsEnd { get; set; } public int? UpdatePlaylistsUniqueTracks { get; set; } public DateTime? PlaylistFollowersStart { get; set; } public DateTime? PlaylistFollowersEnd { get; set; } public DateTime? ReleaseStart { get; set; } public DateTime? ReleaseEnd { get; set; } public int? UniqueAlbumReleases { get; set; } public int? UniqueSingleReleases { get; set; } public DateTime? ImportBuzzPlaylistsStart { get; set; } public DateTime? ImportBuzzPlaylistsEnd { get; set; } public int? ImportBuzzPlaylistsNumUsers { get; set; } public int? ImportBuzzPlaylistsNumNewPlaylists { get; set; } public DateTime? ImportPlaylistsStart { get; set; } public DateTime? ImportPlaylistsEnd { get; set; } public int? ImportPlaylistsNumImported { get; set; } public DateTime? PromotedTracksStart { get; set; } public DateTime? PromotedTracksEnd { get; set; } public int? PromotedTracksUniqueTracks { get; set; } public DateTime? SpotifyAnalyticsStart { get; set; } public DateTime? SpotifyAnalyticsEnd { get; set; } public DateTime? SpotifyAnalyticsDateImported { get; set; } public int? SpotifyAnalyticsRowsImported { get; set; } public DateTime? SpotifyChartsStart { get; set; } public DateTime? SpotifyChartsEnd { get; set; } public int? SpotifyChartsRowsImported { get; set; } public DateTime? StagingDBSyncStart { get; set; } public DateTime? StagingDBSyncEnd { get; set; } } }