diff --git a/Hangfire.Oracle/Configuration/HangfireConfiguration.cs b/Hangfire.Oracle/Configuration/HangfireConfiguration.cs new file mode 100644 index 0000000..2c06240 --- /dev/null +++ b/Hangfire.Oracle/Configuration/HangfireConfiguration.cs @@ -0,0 +1,12 @@ +using System.Collections.Generic; + +namespace Hangfire.Oracle.Core.Configuration +{ + public class HangfireConfiguration + { + public string SchemaName { get; set; } = string.Empty; + public string InstanceName { get; set; } = string.Empty; + public SequenceConfiguration Sequence { get; set; } = new SequenceConfiguration(); + public Dictionary Tables { get; set; } = new Dictionary(); + } +} \ No newline at end of file diff --git a/Hangfire.Oracle/Configuration/HangfireConfigurationLoader.cs b/Hangfire.Oracle/Configuration/HangfireConfigurationLoader.cs new file mode 100644 index 0000000..943e181 --- /dev/null +++ b/Hangfire.Oracle/Configuration/HangfireConfigurationLoader.cs @@ -0,0 +1,71 @@ +using System; +using System.IO; +using Newtonsoft.Json; + +namespace Hangfire.Oracle.Core.Configuration +{ + public static class HangfireConfigurationLoader + { + public static HangfireConfiguration LoadFromJson(string jsonFilePath) + { + if (string.IsNullOrWhiteSpace(jsonFilePath)) + { + throw new ArgumentException("JSON file path cannot be null or empty.", nameof(jsonFilePath)); + } + + if (!File.Exists(jsonFilePath)) + { + throw new FileNotFoundException($"Configuration file not found: {jsonFilePath}"); + } + + var jsonContent = File.ReadAllText(jsonFilePath); + return DeserializeFromJson(jsonContent); + } + + public static HangfireConfiguration DeserializeFromJson(string jsonContent) + { + if (string.IsNullOrWhiteSpace(jsonContent)) + { + throw new ArgumentException("JSON content cannot be null or empty.", nameof(jsonContent)); + } + + try + { + var settings = new JsonSerializerSettings + { + MissingMemberHandling = MissingMemberHandling.Ignore, + NullValueHandling = NullValueHandling.Ignore + }; + + var config = JsonConvert.DeserializeObject(jsonContent, settings); + + if (config == null) + { + throw new InvalidOperationException("Failed to deserialize Hangfire configuration."); + } + + ValidateAndInitialize(config); + + return config; + } + catch (JsonException ex) + { + throw new InvalidOperationException("Invalid JSON format for Hangfire configuration.", ex); + } + } + + private static void ValidateAndInitialize(HangfireConfiguration config) + { + config.Tables = config.Tables ?? new System.Collections.Generic.Dictionary(); + config.Sequence = config.Sequence ?? new SequenceConfiguration(); + + foreach (var table in config.Tables) + { + if (string.IsNullOrWhiteSpace(table.Value)) + { + throw new InvalidOperationException($"Table name for '{table.Key}' cannot be null or empty."); + } + } + } + } +} \ No newline at end of file diff --git a/Hangfire.Oracle/Configuration/HangfireTableNameProvider.cs b/Hangfire.Oracle/Configuration/HangfireTableNameProvider.cs new file mode 100644 index 0000000..4065340 --- /dev/null +++ b/Hangfire.Oracle/Configuration/HangfireTableNameProvider.cs @@ -0,0 +1,88 @@ +using System.Collections.Generic; + +namespace Hangfire.Oracle.Core.Configuration +{ + public class HangfireTableNameProvider + { + private readonly HangfireConfiguration _config; + private readonly Dictionary _defaultTableNames; + private readonly string _instancePrefix; + + public HangfireTableNameProvider(HangfireConfiguration config) + { + _config = config ?? new HangfireConfiguration(); + + _defaultTableNames = new Dictionary + { + { "Job", "HF_JOB" }, + { "JobParameter", "HF_JOB_PARAMETER" }, + { "JobQueue", "HF_JOB_QUEUE" }, + { "JobState", "HF_JOB_STATE" }, + { "Server", "HF_SERVER" }, + { "Set", "HF_SET" }, + { "List", "HF_LIST" }, + { "Hash", "HF_HASH" }, + { "Counter", "HF_COUNTER" }, + { "AggregatedCounter", "HF_AGGREGATED_COUNTER" }, + { "DistributedLock", "HF_DISTRIBUTED_LOCK" } + }; + + _instancePrefix = DetermineInstancePrefix(); + } + + private string DetermineInstancePrefix() + { + if (!string.IsNullOrWhiteSpace(_config.InstanceName)) + { + return _config.InstanceName.ToUpper(); + } + return "HF"; + } + + public string GetTableName(string logicalName) + { + if (_config.Tables != null && _config.Tables.TryGetValue(logicalName, out var customName)) + { + return customName; + } + return _defaultTableNames.TryGetValue(logicalName, out var defaultName) ? defaultName : logicalName; + } + + public string GetSchemaName() + { + return _config.SchemaName ?? string.Empty; + } + + public string GetPrimarySequenceName() + { + return _config.Sequence?.PrimarySequenceName ?? "HF_SEQUENCE"; + } + + public string GetJobIdSequenceName() + { + return _config.Sequence?.JobIdSequenceName ?? "HF_JOB_ID_SEQ"; + } + + public string GetInstancePrefix() + { + return _instancePrefix; + } + + public string GetConstraintName(string tableName, string constraintType) + { + return $"{_instancePrefix}_{constraintType}_{ShortenTableName(tableName)}"; + } + + public string GetIndexName(string tableName, string indexType = "IDX") + { + return $"{_instancePrefix}_{indexType}_{ShortenTableName(tableName)}"; + } + + private string ShortenTableName(string tableName) + { + var name = tableName.Replace($"{_instancePrefix}_", "").Replace("HF_", ""); + + return name.Length > 20 ? name.Substring(0, 20) : name; + } + } +} \ No newline at end of file diff --git a/Hangfire.Oracle/Configuration/SequenceConfiguration.cs b/Hangfire.Oracle/Configuration/SequenceConfiguration.cs new file mode 100644 index 0000000..97a57bb --- /dev/null +++ b/Hangfire.Oracle/Configuration/SequenceConfiguration.cs @@ -0,0 +1,8 @@ +namespace Hangfire.Oracle.Core.Configuration +{ + public class SequenceConfiguration + { + public string PrimarySequenceName { get; set; } = "HF_SEQUENCE"; + public string JobIdSequenceName { get; set; } = "HF_JOB_ID_SEQ"; + } +} \ No newline at end of file diff --git a/Hangfire.Oracle/CountersAggregator.cs b/Hangfire.Oracle/CountersAggregator.cs index 11d59b3..8d8a5df 100644 --- a/Hangfire.Oracle/CountersAggregator.cs +++ b/Hangfire.Oracle/CountersAggregator.cs @@ -1,8 +1,6 @@ using System; using System.Threading; - using Dapper; - using Hangfire.Logging; using Hangfire.Server; @@ -26,6 +24,9 @@ public CountersAggregator(OracleStorage storage, TimeSpan interval) _interval = interval; } + private string T(string logicalName) => _storage.TableNameProvider.GetTableName(logicalName); + private string GetPrimarySequence() => _storage.TableNameProvider.GetPrimarySequenceName(); + public void Execute(CancellationToken cancellationToken) { Logger.DebugFormat("Aggregating records in 'Counter' table..."); @@ -54,35 +55,26 @@ public override string ToString() return GetType().ToString(); } - private static string GetMergeQuery() + private string GetMergeQuery() { - return @" -BEGIN - MERGE INTO HF_AGGREGATED_COUNTER AC - USING ( SELECT KEY, SUM (VALUE) AS VALUE, MAX (EXPIRE_AT) AS EXPIRE_AT - FROM (SELECT KEY, VALUE, EXPIRE_AT - FROM HF_COUNTER - WHERE ROWNUM <= :COUNT) TMP - GROUP BY KEY) C - ON (AC.KEY = C.KEY) - WHEN MATCHED - THEN - UPDATE SET VALUE = AC.VALUE + C.VALUE, EXPIRE_AT = GREATEST (EXPIRE_AT, C.EXPIRE_AT) - WHEN NOT MATCHED - THEN - INSERT (ID - ,KEY - ,VALUE - ,EXPIRE_AT) - VALUES (HF_SEQUENCE.NEXTVAL - ,C.KEY - ,C.VALUE - ,C.EXPIRE_AT); + return $@" + BEGIN + MERGE INTO {T("AggregatedCounter")} AC + USING (SELECT KEY, SUM(VALUE) AS VALUE, MAX(EXPIRE_AT) AS EXPIRE_AT + FROM (SELECT KEY, VALUE, EXPIRE_AT + FROM {T("Counter")} + WHERE ROWNUM <= :COUNT) TMP + GROUP BY KEY) C + ON (AC.KEY = C.KEY) + WHEN MATCHED THEN + UPDATE SET VALUE = AC.VALUE + C.VALUE, EXPIRE_AT = GREATEST(EXPIRE_AT, C.EXPIRE_AT) + WHEN NOT MATCHED THEN + INSERT (ID, KEY, VALUE, EXPIRE_AT) + VALUES ({GetPrimarySequence()}.NEXTVAL, C.KEY, C.VALUE, C.EXPIRE_AT); - DELETE FROM HF_COUNTER - WHERE ROWNUM <= :COUNT; -END; -"; + DELETE FROM {T("Counter")} + WHERE ROWNUM <= :COUNT; + END;"; } } } diff --git a/Hangfire.Oracle/Entities/EntityUtils.cs b/Hangfire.Oracle/Entities/EntityUtils.cs index 50e1532..7af775d 100644 --- a/Hangfire.Oracle/Entities/EntityUtils.cs +++ b/Hangfire.Oracle/Entities/EntityUtils.cs @@ -1,19 +1,18 @@ using System.Data; - using Dapper; namespace Hangfire.Oracle.Core.Entities { public static class EntityUtils { - public static long GetNextId(this IDbConnection connection) + public static long GetNextId(this IDbConnection connection, string sequenceName = "HF_SEQUENCE") { - return connection.QuerySingle("SELECT HF_SEQUENCE.NEXTVAL FROM dual"); + return connection.QuerySingle($"SELECT {sequenceName}.NEXTVAL FROM dual"); } - public static long GetNextJobId(this IDbConnection connection) + public static long GetNextJobId(this IDbConnection connection, string sequenceName = "HF_JOB_ID_SEQ") { - return connection.QuerySingle("SELECT HF_JOB_ID_SEQ.NEXTVAL FROM dual"); + return connection.QuerySingle($"SELECT {sequenceName}.NEXTVAL FROM dual"); } } } diff --git a/Hangfire.Oracle/ExpirationManager.cs b/Hangfire.Oracle/ExpirationManager.cs index d9e86f0..b7d1f61 100644 --- a/Hangfire.Oracle/ExpirationManager.cs +++ b/Hangfire.Oracle/ExpirationManager.cs @@ -2,9 +2,7 @@ using System.Collections.Generic; using System.Data.Common; using System.Threading; - using Dapper; - using Hangfire.Logging; using Hangfire.Server; @@ -21,19 +19,6 @@ internal class ExpirationManager : IServerComponent private static readonly TimeSpan DelayBetweenPasses = TimeSpan.FromSeconds(1); private const int NumberOfRecordsInSinglePass = 1000; - private static readonly List> TablesToProcess = new List> - { - // This list must be sorted in dependency order - new Tuple("HF_JOB_PARAMETER", true), - new Tuple("HF_JOB_QUEUE", true), - new Tuple("HF_JOB_STATE", true), - new Tuple("HF_AGGREGATED_COUNTER", false), - new Tuple("HF_LIST", false), - new Tuple("HF_SET", false), - new Tuple("HF_HASH", false), - new Tuple("HF_JOB", false) - }; - private readonly OracleStorage _storage; private readonly TimeSpan _checkInterval; @@ -48,9 +33,24 @@ public ExpirationManager(OracleStorage storage, TimeSpan checkInterval) _checkInterval = checkInterval; } + private string T(string logicalName) => _storage.TableNameProvider.GetTableName(logicalName); + public void Execute(CancellationToken cancellationToken) { - foreach (var tuple in TablesToProcess) + // This list must be sorted in dependency order + var tablesToProcess = new List> + { + new Tuple("JobParameter", true), + new Tuple("JobQueue", true), + new Tuple("JobState", true), + new Tuple("AggregatedCounter", false), + new Tuple("List", false), + new Tuple("Set", false), + new Tuple("Hash", false), + new Tuple("Job", false) + }; + + foreach (var tuple in tablesToProcess) { Logger.DebugFormat("Removing outdated records from table '{0}'...", tuple.Item1); @@ -66,10 +66,10 @@ public void Execute(CancellationToken cancellationToken) using (new OracleDistributedLock(connection, DistributedLockKey, DefaultLockTimeout, cancellationToken).Acquire()) { - var query = $"DELETE FROM {tuple.Item1} WHERE EXPIRE_AT < :NOW AND ROWNUM <= :COUNT"; + var query = $"DELETE FROM {T(tuple.Item1)} WHERE EXPIRE_AT < :NOW AND ROWNUM <= :COUNT"; if (tuple.Item2) { - query = $"DELETE FROM {tuple.Item1} WHERE JOB_ID IN (SELECT ID FROM HF_JOB WHERE EXPIRE_AT < :NOW AND ROWNUM <= :COUNT)"; + query = $"DELETE FROM {T(tuple.Item1)} WHERE JOB_ID IN (SELECT ID FROM {T("Job")} WHERE EXPIRE_AT < :NOW AND ROWNUM <= :COUNT)"; } removedCount = connection.Execute(query, new { NOW = DateTime.UtcNow, COUNT = NumberOfRecordsInSinglePass }); } diff --git a/Hangfire.Oracle/Hangfire.Oracle.Core.csproj b/Hangfire.Oracle/Hangfire.Oracle.Core.csproj index a7f1774..13631ae 100644 --- a/Hangfire.Oracle/Hangfire.Oracle.Core.csproj +++ b/Hangfire.Oracle/Hangfire.Oracle.Core.csproj @@ -20,12 +20,6 @@ Fix missing prefix in merge query 1.3.1 - - - - - - diff --git a/Hangfire.Oracle/Install.sql b/Hangfire.Oracle/Install.sql deleted file mode 100644 index 73f145f..0000000 --- a/Hangfire.Oracle/Install.sql +++ /dev/null @@ -1,282 +0,0 @@ -CREATE SEQUENCE HF_SEQUENCE START WITH 1 MAXVALUE 9999999999999999999999999999 MINVALUE 1 NOCYCLE CACHE 20 NOORDER; - -CREATE SEQUENCE HF_JOB_ID_SEQ START WITH 1 MAXVALUE 9999999999999999999999999999 MINVALUE 1 NOCYCLE CACHE 20 NOORDER; - --- ---------------------------- --- Table structure for `Job` --- ---------------------------- - -CREATE TABLE HF_JOB (ID NUMBER (10) - ,STATE_ID NUMBER (10) - ,STATE_NAME NVARCHAR2 (20) - ,INVOCATION_DATA NCLOB - ,ARGUMENTS NCLOB - ,CREATED_AT TIMESTAMP (4) - ,EXPIRE_AT TIMESTAMP (4)) -LOB (INVOCATION_DATA) STORE AS BASICFILE - (ENABLE STORAGE IN ROW - CHUNK 8192 - RETENTION - NOCACHE LOGGING) -LOB (ARGUMENTS) STORE AS BASICFILE - (ENABLE STORAGE IN ROW - CHUNK 8192 - RETENTION - NOCACHE LOGGING) -LOGGING -NOCOMPRESS -NOCACHE -NOPARALLEL -MONITORING; - -ALTER TABLE HF_JOB ADD ( - PRIMARY KEY - (ID) - USING INDEX - ENABLE VALIDATE); - - --- ---------------------------- --- Table structure for `Counter` --- ---------------------------- - -CREATE TABLE HF_COUNTER (ID NUMBER (10) - ,KEY NVARCHAR2 (255) - ,VALUE NUMBER (10) - ,EXPIRE_AT TIMESTAMP (4)) -LOGGING -NOCOMPRESS -NOCACHE -NOPARALLEL -MONITORING; - -ALTER TABLE HF_COUNTER ADD ( - PRIMARY KEY - (ID) - USING INDEX - ENABLE VALIDATE); - - -CREATE TABLE HF_AGGREGATED_COUNTER (ID NUMBER (10) - ,KEY NVARCHAR2 (255) - ,VALUE NUMBER (10) - ,EXPIRE_AT TIMESTAMP (4)) -LOGGING -NOCOMPRESS -NOCACHE -NOPARALLEL -MONITORING; - -ALTER TABLE HF_AGGREGATED_COUNTER ADD ( - PRIMARY KEY - (ID) - USING INDEX - ENABLE VALIDATE, - UNIQUE (KEY) - USING INDEX - ENABLE VALIDATE); - - --- ---------------------------- --- Table structure for `DistributedLock` --- ---------------------------- - -CREATE TABLE HF_DISTRIBUTED_LOCK ("RESOURCE" NVARCHAR2 (100), CREATED_AT TIMESTAMP (4)) -LOGGING -NOCOMPRESS -NOCACHE -NOPARALLEL -MONITORING; - - --- ---------------------------- --- Table structure for `Hash` --- ---------------------------- - -CREATE TABLE HF_HASH (ID NUMBER (10) - ,KEY NVARCHAR2 (255) - ,VALUE NCLOB - ,EXPIRE_AT TIMESTAMP (4) - ,FIELD NVARCHAR2 (40)) -LOB (VALUE) STORE AS BASICFILE - (ENABLE STORAGE IN ROW - CHUNK 8192 - RETENTION - NOCACHE LOGGING) -LOGGING -NOCOMPRESS -NOCACHE -NOPARALLEL -MONITORING; - -ALTER TABLE HF_HASH ADD ( - PRIMARY KEY - (ID) - USING INDEX - ENABLE VALIDATE, - UNIQUE (KEY, FIELD) - USING INDEX - ENABLE VALIDATE); - - --- ---------------------------- --- Table structure for `JobParameter` --- ---------------------------- - -CREATE TABLE HF_JOB_PARAMETER (ID NUMBER (10) - ,NAME NVARCHAR2 (40) - ,VALUE NCLOB - ,JOB_ID NUMBER (10)) -LOB (VALUE) STORE AS BASICFILE - (ENABLE STORAGE IN ROW - CHUNK 8192 - RETENTION - NOCACHE LOGGING) -LOGGING -NOCOMPRESS -NOCACHE -NOPARALLEL -MONITORING; - -ALTER TABLE HF_JOB_PARAMETER ADD ( - PRIMARY KEY - (ID) - USING INDEX - ENABLE VALIDATE); - -ALTER TABLE HF_JOB_PARAMETER ADD ( - CONSTRAINT FK_JOB_PARAMETER_JOB - FOREIGN KEY (JOB_ID) - REFERENCES HF_JOB (ID) - ON DELETE CASCADE ENABLE VALIDATE); - - --- ---------------------------- --- Table structure for `JobQueue` --- ---------------------------- - -CREATE TABLE HF_JOB_QUEUE (ID NUMBER (10) - ,JOB_ID NUMBER (10) - ,QUEUE NVARCHAR2 (50) - ,FETCHED_AT TIMESTAMP (4) - ,FETCH_TOKEN NVARCHAR2 (36)) -LOGGING -NOCOMPRESS -NOCACHE -NOPARALLEL -MONITORING; - -ALTER TABLE HF_JOB_QUEUE ADD ( - PRIMARY KEY - (ID) - USING INDEX - ENABLE VALIDATE); - -ALTER TABLE HF_JOB_QUEUE ADD ( - CONSTRAINT FK_JOB_QUEUE_JOB - FOREIGN KEY (JOB_ID) - REFERENCES HF_JOB (ID) - ON DELETE CASCADE ENABLE VALIDATE); - - --- ---------------------------- --- Table structure for `JobState` --- ---------------------------- - -CREATE TABLE HF_JOB_STATE (ID NUMBER (10) - ,JOB_ID NUMBER (10) - ,NAME NVARCHAR2 (20) - ,REASON NVARCHAR2 (100) - ,CREATED_AT TIMESTAMP (4) - ,DATA NCLOB) -LOB (DATA) STORE AS BASICFILE - (ENABLE STORAGE IN ROW - CHUNK 8192 - RETENTION - NOCACHE LOGGING) -LOGGING -NOCOMPRESS -NOCACHE -NOPARALLEL -MONITORING; - -ALTER TABLE HF_JOB_STATE ADD ( - PRIMARY KEY - (ID) - USING INDEX - ENABLE VALIDATE); - -ALTER TABLE HF_JOB_STATE ADD ( - CONSTRAINT FK_JOB_STATE_JOB - FOREIGN KEY (JOB_ID) - REFERENCES HF_JOB (ID) - ON DELETE CASCADE ENABLE VALIDATE); - - --- ---------------------------- --- Table structure for `Server` --- ---------------------------- - -CREATE TABLE HF_SERVER (ID NVARCHAR2 (100), DATA NCLOB, LAST_HEART_BEAT TIMESTAMP (4)) -LOB (DATA) STORE AS BASICFILE - (ENABLE STORAGE IN ROW - CHUNK 8192 - RETENTION - NOCACHE LOGGING) -LOGGING -NOCOMPRESS -NOCACHE -NOPARALLEL -MONITORING; - -ALTER TABLE HF_SERVER ADD ( - PRIMARY KEY - (ID) - USING INDEX - ENABLE VALIDATE); - - --- ---------------------------- --- Table structure for `Set` --- ---------------------------- - -CREATE TABLE HF_SET (ID NUMBER (10) - ,KEY NVARCHAR2 (255) - ,VALUE NVARCHAR2 (255) - ,SCORE FLOAT (126) - ,EXPIRE_AT TIMESTAMP (4)) -LOGGING -NOCOMPRESS -NOCACHE -NOPARALLEL -MONITORING; - -ALTER TABLE HF_SET ADD ( - PRIMARY KEY - (ID) - USING INDEX - ENABLE VALIDATE, - UNIQUE (KEY, VALUE) - USING INDEX - ENABLE VALIDATE); - -CREATE TABLE HF_LIST (ID NUMBER (10) - ,KEY NVARCHAR2 (255) - ,VALUE NCLOB - ,EXPIRE_AT TIMESTAMP (4)) -LOB (VALUE) STORE AS BASICFILE - (ENABLE STORAGE IN ROW - CHUNK 8192 - RETENTION - NOCACHE LOGGING) -LOGGING -NOCOMPRESS -NOCACHE -NOPARALLEL -MONITORING; - -ALTER TABLE HF_LIST ADD ( - PRIMARY KEY - (ID) - USING INDEX - ENABLE VALIDATE); \ No newline at end of file diff --git a/Hangfire.Oracle/JobQueue/IPersistentJobQueue.cs b/Hangfire.Oracle/JobQueue/IPersistentJobQueue.cs index b8617bf..d36d9b7 100644 --- a/Hangfire.Oracle/JobQueue/IPersistentJobQueue.cs +++ b/Hangfire.Oracle/JobQueue/IPersistentJobQueue.cs @@ -1,6 +1,5 @@ using System.Data; using System.Threading; - using Hangfire.Storage; namespace Hangfire.Oracle.Core.JobQueue diff --git a/Hangfire.Oracle/JobQueue/OracleFetchedJob.cs b/Hangfire.Oracle/JobQueue/OracleFetchedJob.cs index e294db7..fe99eda 100644 --- a/Hangfire.Oracle/JobQueue/OracleFetchedJob.cs +++ b/Hangfire.Oracle/JobQueue/OracleFetchedJob.cs @@ -1,9 +1,7 @@ using System; using System.Data; using System.Globalization; - using Dapper; - using Hangfire.Logging; using Hangfire.Storage; @@ -34,9 +32,10 @@ public OracleFetchedJob(OracleStorage storage, IDbConnection connection, Fetched Queue = fetchedJob.Queue; } + private string T(string logicalName) => _storage.TableNameProvider.GetTableName(logicalName); + public void Dispose() { - if (_disposed) return; if (!_removedFromQueue && !_requeued) @@ -53,10 +52,8 @@ public void RemoveFromQueue() { Logger.TraceFormat("RemoveFromQueue JobId={0}", JobId); - //todo: unit test _connection.Execute( - "DELETE FROM HF_JOB_QUEUE " + - " WHERE ID = :ID", + $"DELETE FROM {T("JobQueue")} WHERE ID = :ID", new { ID = _id @@ -69,11 +66,8 @@ public void Requeue() { Logger.TraceFormat("Requeue JobId={0}", JobId); - //todo: unit test _connection.Execute( - "UPDATE HF_JOB_QUEUE " + - " SET FETCHED_AT = null " + - " WHERE ID = :ID", + $"UPDATE {T("JobQueue")} SET FETCHED_AT = null WHERE ID = :ID", new { ID = _id diff --git a/Hangfire.Oracle/JobQueue/OracleJobQueue.cs b/Hangfire.Oracle/JobQueue/OracleJobQueue.cs index c401459..f901ef8 100644 --- a/Hangfire.Oracle/JobQueue/OracleJobQueue.cs +++ b/Hangfire.Oracle/JobQueue/OracleJobQueue.cs @@ -2,9 +2,7 @@ using System.Data; using System.Data.Common; using System.Threading; - using Dapper; - using Hangfire.Logging; using Hangfire.Storage; @@ -16,12 +14,17 @@ internal class OracleJobQueue : IPersistentJobQueue private readonly OracleStorage _storage; private readonly OracleStorageOptions _options; + + private string GetPrimarySequence() => _storage.TableNameProvider.GetPrimarySequenceName(); + public OracleJobQueue(OracleStorage storage, OracleStorageOptions options) { _storage = storage ?? throw new ArgumentNullException(nameof(storage)); _options = options ?? throw new ArgumentNullException(nameof(options)); } + private string T(string logicalName) => _storage.TableNameProvider.GetTableName(logicalName); + public IFetchedJob Dequeue(string[] queues, CancellationToken cancellationToken) { if (queues == null) throw new ArgumentNullException(nameof(queues)); @@ -41,11 +44,11 @@ public IFetchedJob Dequeue(string[] queues, CancellationToken cancellationToken) { var token = Guid.NewGuid().ToString(); - var nUpdated = connection.Execute(@" -UPDATE HF_JOB_QUEUE - SET FETCHED_AT = SYS_EXTRACT_UTC (SYSTIMESTAMP), FETCH_TOKEN = :FETCH_TOKEN - WHERE (FETCHED_AT IS NULL OR FETCHED_AT < SYS_EXTRACT_UTC (SYSTIMESTAMP) + numToDSInterval(:TIMEOUT, 'second' )) AND (QUEUE IN :QUEUES) AND ROWNUM = 1 -", + var nUpdated = connection.Execute($@" + UPDATE {T("JobQueue")} + SET FETCHED_AT = SYS_EXTRACT_UTC(SYSTIMESTAMP), FETCH_TOKEN = :FETCH_TOKEN + WHERE (FETCHED_AT IS NULL OR FETCHED_AT < SYS_EXTRACT_UTC(SYSTIMESTAMP) + numToDSInterval(:TIMEOUT, 'second')) + AND (QUEUE IN :QUEUES) AND ROWNUM = 1", new { QUEUES = queues, @@ -55,18 +58,14 @@ UPDATE HF_JOB_QUEUE if (nUpdated != 0) { - fetchedJob = - connection - .QuerySingle( - @" - SELECT ID as Id, JOB_ID as JobId, QUEUE as Queue - FROM HF_JOB_QUEUE - WHERE FETCH_TOKEN = :FETCH_TOKEN -", - new - { - FETCH_TOKEN = token - }); + fetchedJob = connection.QuerySingle( + $@"SELECT ID as Id, JOB_ID as JobId, QUEUE as Queue + FROM {T("JobQueue")} + WHERE FETCH_TOKEN = :FETCH_TOKEN", + new + { + FETCH_TOKEN = token + }); } } } @@ -92,7 +91,9 @@ FROM HF_JOB_QUEUE public void Enqueue(IDbConnection connection, string queue, string jobId) { Logger.TraceFormat("Enqueue JobId={0} Queue={1}", jobId, queue); - connection.Execute("INSERT INTO HF_JOB_QUEUE (ID, JOB_ID, QUEUE) VALUES (HF_SEQUENCE.NEXTVAL, :JOB_ID, :QUEUE)", new { JOB_ID = jobId, QUEUE = queue }); + connection.Execute( + $"INSERT INTO {T("JobQueue")} (ID, JOB_ID, QUEUE) VALUES ({GetPrimarySequence()}.NEXTVAL, :JOB_ID, :QUEUE)", + new { JOB_ID = jobId, QUEUE = queue }); } } } \ No newline at end of file diff --git a/Hangfire.Oracle/JobQueue/OracleJobQueueMonitoringApi.cs b/Hangfire.Oracle/JobQueue/OracleJobQueueMonitoringApi.cs index 79a3bae..afc1849 100644 --- a/Hangfire.Oracle/JobQueue/OracleJobQueueMonitoringApi.cs +++ b/Hangfire.Oracle/JobQueue/OracleJobQueueMonitoringApi.cs @@ -1,7 +1,6 @@ using System; using System.Collections.Generic; using System.Linq; - using Dapper; namespace Hangfire.Oracle.Core.JobQueue @@ -14,11 +13,14 @@ internal class OracleJobQueueMonitoringApi : IPersistentJobQueueMonitoringApi private DateTime _cacheUpdated; private readonly OracleStorage _storage; + public OracleJobQueueMonitoringApi(OracleStorage storage) { _storage = storage ?? throw new ArgumentNullException(nameof(storage)); } + private string T(string logicalName) => _storage.TableNameProvider.GetTableName(logicalName); + public IEnumerable GetQueues() { lock (_cacheLock) @@ -27,7 +29,8 @@ public IEnumerable GetQueues() { var result = _storage.UseConnection(connection => { - return connection.Query("SELECT DISTINCT(QUEUE) as QUEUE FROM HF_JOB_QUEUE").Select(x => (string)x.QUEUE).ToList(); + return connection.Query($"SELECT DISTINCT(QUEUE) as QUEUE FROM {T("JobQueue")}") + .Select(x => (string)x.QUEUE).ToList(); }); _queuesCache = result; @@ -40,13 +43,12 @@ public IEnumerable GetQueues() public IEnumerable GetEnqueuedJobIds(string queue, int from, int perPage) { - const string sqlQuery = @" -SELECT JOB_ID AS JobId - FROM (SELECT JOB_ID, RANK () OVER (ORDER BY ID) AS RANK - FROM HF_JOB_QUEUE - WHERE QUEUE = :QUEUE) - WHERE RANK BETWEEN :S AND :E -"; + var sqlQuery = $@" + SELECT JOB_ID AS JobId + FROM (SELECT JOB_ID, RANK() OVER (ORDER BY ID) AS RANK + FROM {T("JobQueue")} + WHERE QUEUE = :QUEUE AND FETCHED_AT IS NULL) + WHERE RANK BETWEEN :S AND :E"; return _storage.UseConnection(connection => connection.Query(sqlQuery, new { QUEUE = queue, S = from + 1, E = from + perPage })); @@ -54,18 +56,33 @@ FROM HF_JOB_QUEUE public IEnumerable GetFetchedJobIds(string queue, int from, int perPage) { - return Enumerable.Empty(); + var sqlQuery = $@" + SELECT JOB_ID AS JobId + FROM (SELECT JOB_ID, RANK() OVER (ORDER BY ID) AS RANK + FROM {T("JobQueue")} + WHERE QUEUE = :QUEUE AND FETCHED_AT IS NOT NULL) + WHERE RANK BETWEEN :S AND :E"; + + return _storage.UseConnection(connection => + connection.Query(sqlQuery, new { QUEUE = queue, S = from + 1, E = from + perPage })); } public EnqueuedAndFetchedCountDto GetEnqueuedAndFetchedCount(string queue) { return _storage.UseConnection(connection => { - var result = connection.QuerySingle("SELECT COUNT(ID) FROM HF_JOB_QUEUE WHERE QUEUE = :QUEUE", new { QUEUE = queue }); + var enqueuedCount = connection.QuerySingle( + $"SELECT COUNT(ID) FROM {T("JobQueue")} WHERE QUEUE = :QUEUE AND FETCHED_AT IS NULL", + new { QUEUE = queue }); + + var fetchedCount = connection.QuerySingle( + $"SELECT COUNT(ID) FROM {T("JobQueue")} WHERE QUEUE = :QUEUE AND FETCHED_AT IS NOT NULL", + new { QUEUE = queue }); return new EnqueuedAndFetchedCountDto { - EnqueuedCount = result + EnqueuedCount = enqueuedCount, + FetchedCount = fetchedCount }; }); } diff --git a/Hangfire.Oracle/Monitoring/OracleMonitoringApi.cs b/Hangfire.Oracle/Monitoring/OracleMonitoringApi.cs index ecf7155..13befad 100644 --- a/Hangfire.Oracle/Monitoring/OracleMonitoringApi.cs +++ b/Hangfire.Oracle/Monitoring/OracleMonitoringApi.cs @@ -2,9 +2,7 @@ using System.Collections.Generic; using System.Data; using System.Linq; - using Dapper; - using Hangfire.Annotations; using Hangfire.Common; using Hangfire.Oracle.Core.Entities; @@ -26,6 +24,8 @@ public OracleMonitoringApi([NotNull] OracleStorage storage, int? jobListLimit) _jobListLimit = jobListLimit; } + private string T(string logicalName) => _storage.TableNameProvider.GetTableName(logicalName); + public IList Queues() { var tuples = _storage.QueueProviders @@ -59,8 +59,8 @@ public IList Servers() { return UseConnection>(connection => { - var servers = - connection.Query("SELECT ID as ID, DATA as Data, LAST_HEART_BEAT as LastHeartBeat FROM HF_SERVER").ToList(); + var servers = connection.Query( + $"SELECT ID as ID, DATA as Data, LAST_HEART_BEAT as LastHeartBeat FROM {T("Server")}").ToList(); var result = new List(); @@ -85,45 +85,44 @@ public JobDetailsDto JobDetails(string jobId) { return UseConnection(connection => { - const string jobQuery = @" - SELECT ID AS Id - ,STATE_ID AS StateId - ,STATE_NAME AS StateName - ,INVOCATION_DATA AS InvocationData - ,ARGUMENTS AS Arguments - ,CREATED_AT AS CreatedAt - ,EXPIRE_AT AS ExpireAt - FROM HF_JOB - WHERE ID = :ID -"; + var jobQuery = $@" + SELECT ID AS Id + ,STATE_ID AS StateId + ,STATE_NAME AS StateName + ,INVOCATION_DATA AS InvocationData + ,ARGUMENTS AS Arguments + ,CREATED_AT AS CreatedAt + ,EXPIRE_AT AS ExpireAt + FROM {T("Job")} + WHERE ID = :ID"; + var sqlJob = connection.QuerySingleOrDefault(jobQuery, new { ID = jobId }); if (sqlJob == null) { return null; } - const string jobParametersQuery = @" - SELECT ID AS Id - ,NAME AS Name - ,VALUE AS Value - ,JOB_ID AS JobId - FROM HF_JOB_PARAMETER - WHERE JOB_ID = :ID -"; - - var jobParameters = connection.Query(jobParametersQuery, new { ID = jobId }).ToDictionary(x => x.Name, parameter => parameter.Value); - - const string jobStatesQuery = @" - SELECT ID AS Id - ,JOB_ID AS JobId - ,NAME AS Name - ,REASON AS Reason - ,CREATED_AT AS CreatedAt - ,DATA AS Data - FROM HF_JOB_STATE - WHERE JOB_ID = :ID - ORDER BY ID DESC -"; + var jobParametersQuery = $@" + SELECT ID AS Id + ,NAME AS Name + ,VALUE AS Value + ,JOB_ID AS JobId + FROM {T("JobParameter")} + WHERE JOB_ID = :ID"; + + var jobParameters = connection.Query(jobParametersQuery, new { ID = jobId }) + .ToDictionary(x => x.Name, parameter => parameter.Value); + + var jobStatesQuery = $@" + SELECT ID AS Id + ,JOB_ID AS JobId + ,NAME AS Name + ,REASON AS Reason + ,CREATED_AT AS CreatedAt + ,DATA AS Data + FROM {T("JobState")} + WHERE JOB_ID = :ID + ORDER BY ID DESC"; var jobStates = connection.Query(jobStatesQuery, new { ID = jobId }).Select(x => new StateHistoryDto { @@ -146,17 +145,16 @@ ORDER BY ID DESC public StatisticsDto GetStatistics() { - const string jobQuery = "SELECT COUNT(ID) FROM HF_JOB WHERE STATE_NAME = :STATE_NAME"; - const string succeededQuery = @" - SELECT SUM (VALUE) - FROM (SELECT SUM (VALUE) AS VALUE - FROM HF_COUNTER - WHERE KEY = :KEY - UNION ALL - SELECT VALUE - FROM HF_AGGREGATED_COUNTER - WHERE KEY = :KEY) -"; + var jobQuery = $"SELECT COUNT(ID) FROM {T("Job")} WHERE STATE_NAME = :STATE_NAME"; + var succeededQuery = $@" + SELECT SUM(VALUE) + FROM (SELECT SUM(VALUE) AS VALUE + FROM {T("Counter")} + WHERE KEY = :KEY + UNION ALL + SELECT VALUE + FROM {T("AggregatedCounter")} + WHERE KEY = :KEY)"; var statistics = UseConnection(connection => @@ -166,10 +164,10 @@ FROM HF_AGGREGATED_COUNTER Failed = connection.ExecuteScalar(jobQuery, new { STATE_NAME = "Failed" }), Processing = connection.ExecuteScalar(jobQuery, new { STATE_NAME = "Processing" }), Scheduled = connection.ExecuteScalar(jobQuery, new { STATE_NAME = "Scheduled" }), - Servers = connection.ExecuteScalar("SELECT COUNT(ID) FROM HF_SERVER"), + Servers = connection.ExecuteScalar($"SELECT COUNT(ID) FROM {T("Server")}"), Succeeded = connection.ExecuteScalar(succeededQuery, new { KEY = "stats:succeeded" }), Deleted = connection.ExecuteScalar(succeededQuery, new { KEY = "stats:deleted" }), - Recurring = connection.ExecuteScalar("SELECT COUNT(*) FROM HF_SET WHERE KEY = 'recurring-jobs'") + Recurring = connection.ExecuteScalar($"SELECT COUNT(*) FROM {T("Set")} WHERE KEY = 'recurring-jobs'") }); statistics.Queues = _storage.QueueProviders @@ -228,7 +226,7 @@ public JobList SucceededJobs(int from, int count) Job = job, Result = stateData.ContainsKey("Result") ? stateData["Result"] : null, TotalDuration = stateData.ContainsKey("PerformanceDuration") && stateData.ContainsKey("Latency") - ? (long?)long.Parse(stateData["PerformanceDuration"]) + (long?)long.Parse(stateData["Latency"]) + ? long.Parse(stateData["PerformanceDuration"]) + (long?)long.Parse(stateData["Latency"]) : null, SucceededAt = JobHelper.DeserializeNullableDateTime(stateData["SucceededAt"]) })); @@ -327,8 +325,8 @@ private T UseConnection(Func action) private long GetNumberOfJobsByStateName(IDbConnection connection, string stateName) { var sqlQuery = _jobListLimit.HasValue - ? "SELECT COUNT(ID) FROM (SELECT ID FROM HF_JOB WHERE STATE_NAME = :STATE_NAME AND ROWNUM <= :LIMIT)" - : "SELECT COUNT(ID) FROM HF_JOB WHERE STATE_NAME = :STATE_NAME"; + ? $"SELECT COUNT(ID) FROM (SELECT ID FROM {T("Job")} WHERE STATE_NAME = :STATE_NAME AND ROWNUM <= :LIMIT)" + : $"SELECT COUNT(ID) FROM {T("Job")} WHERE STATE_NAME = :STATE_NAME"; var count = connection.QuerySingle( sqlQuery, @@ -336,6 +334,7 @@ private long GetNumberOfJobsByStateName(IDbConnection connection, string stateNa return count; } + private IPersistentJobQueueMonitoringApi GetQueueApi(string queueName) { var provider = _storage.QueueProviders.GetProvider(queueName); @@ -344,34 +343,33 @@ private IPersistentJobQueueMonitoringApi GetQueueApi(string queueName) return monitoringApi; } - private static JobList GetJobs(IDbConnection connection, int from, int count, string stateName, + private JobList GetJobs(IDbConnection connection, int from, int count, string stateName, Func, TDto> selector) { - const string jobsSql = @" -SELECT Id - ,StateId - ,StateName - ,InvocationData - ,Arguments - ,CreatedAt - ,ExpireAt - ,StateReason - ,StateData - FROM (SELECT J.ID AS Id - ,J.STATE_ID AS StateId - ,J.STATE_NAME AS StateName - ,J.INVOCATION_DATA AS InvocationData - ,J.ARGUMENTS AS Arguments - ,J.CREATED_AT AS CreatedAt - ,J.EXPIRE_AT AS ExpireAt - ,S.REASON AS StateReason - ,S.DATA AS StateData - ,RANK () OVER (ORDER BY J.ID DESC) AS RANK - FROM HF_JOB J LEFT JOIN HF_JOB_STATE S ON J.STATE_ID = S.ID - WHERE J.STATE_NAME = :STATE_NAME - ) - WHERE RANK BETWEEN :S AND :E -"; + var jobsSql = $@" + SELECT Id + ,StateId + ,StateName + ,InvocationData + ,Arguments + ,CreatedAt + ,ExpireAt + ,StateReason + ,StateData + FROM (SELECT J.ID AS Id + ,J.STATE_ID AS StateId + ,J.STATE_NAME AS StateName + ,J.INVOCATION_DATA AS InvocationData + ,J.ARGUMENTS AS Arguments + ,J.CREATED_AT AS CreatedAt + ,J.EXPIRE_AT AS ExpireAt + ,S.REASON AS StateReason + ,S.DATA AS StateData + ,RANK() OVER (ORDER BY J.ID DESC) AS RANK + FROM {T("Job")} J + LEFT JOIN {T("JobState")} S ON J.STATE_ID = S.ID + WHERE J.STATE_NAME = :STATE_NAME) + WHERE RANK BETWEEN :S AND :E"; var jobs = connection.Query(jobsSql, new { STATE_NAME = stateName, S = from + 1, E = from + count }).ToList(); @@ -384,7 +382,7 @@ private static JobList DeserializeJobs(ICollection jobs, Fun foreach (var job in jobs) { - var deserializedData = SerializationHelper.Deserialize>(job.StateData, SerializationOption.User); + var deserializedData = SerializationHelper.Deserialize>(job.StateData, SerializationOption.User); var stateData = deserializedData != null ? new Dictionary(deserializedData, StringComparer.OrdinalIgnoreCase) @@ -413,7 +411,7 @@ private static Job DeserializeJob(string invocationData, string arguments) } } - private static Dictionary GetTimelineStats(IDbConnection connection, string type) + private Dictionary GetTimelineStats(IDbConnection connection, string type) { var endDate = DateTime.UtcNow.Date; var dates = new List(); @@ -428,9 +426,10 @@ private static Dictionary GetTimelineStats(IDbConnection connect return GetTimelineStats(connection, keyMaps); } - private static Dictionary GetTimelineStats(IDbConnection connection, IDictionary keyMaps) + private Dictionary GetTimelineStats(IDbConnection connection, IDictionary keyMaps) { - var valuesMap = connection.Query("SELECT KEY AS Key, VALUE AS Count FROM HF_AGGREGATED_COUNTER WHERE KEY in :KEYS", + var valuesMap = connection.Query( + $"SELECT KEY AS Key, VALUE AS Count FROM {T("AggregatedCounter")} WHERE KEY in :KEYS", new { KEYS = keyMaps.Keys }) .ToDictionary(x => (string)x.KEY, x => (long)x.COUNT); @@ -449,24 +448,25 @@ private static Dictionary GetTimelineStats(IDbConnection connect return result; } - private static JobList EnqueuedJobs(IDbConnection connection, IEnumerable jobIds) + private JobList EnqueuedJobs(IDbConnection connection, IEnumerable jobIds) { var enumerable = jobIds as int[] ?? jobIds.ToArray(); - var enqueuedJobsSql = @" - SELECT J.ID AS Id, J.STATE_ID AS StateId, J.STATE_NAME AS StateName, J.INVOCATION_DATA AS InvocationData, J.ARGUMENTS AS Arguments, J.CREATED_AT AS CreatedAt, J.EXPIRE_AT AS ExpireAt, S.REASON AS StateReason, S.DATA AS StateData - FROM HF_JOB J - LEFT JOIN HF_JOB_STATE S - ON S.ID = J.STATE_ID - WHERE J.ID in :JOB_IDS"; + var enqueuedJobsSql = $@" + SELECT J.ID AS Id, J.STATE_ID AS StateId, J.STATE_NAME AS StateName, J.INVOCATION_DATA AS InvocationData, + J.ARGUMENTS AS Arguments, J.CREATED_AT AS CreatedAt, J.EXPIRE_AT AS ExpireAt, + S.REASON AS StateReason, S.DATA AS StateData + FROM {T("Job")} J + LEFT JOIN {T("JobState")} S ON S.ID = J.STATE_ID + WHERE J.ID in :JOB_IDS"; if (!enumerable.Any()) { - enqueuedJobsSql = @" - SELECT J.ID AS Id, J.STATE_ID AS StateId, J.STATE_NAME AS StateName, J.INVOCATION_DATA AS InvocationData, J.ARGUMENTS AS Arguments, J.CREATED_AT AS CreatedAt, J.EXPIRE_AT AS ExpireAt, S.REASON AS StateReason, S.DATA AS StateData - FROM HF_JOB J - LEFT JOIN HF_JOB_STATE S - ON S.ID = J.STATE_ID -"; + enqueuedJobsSql = $@" + SELECT J.ID AS Id, J.STATE_ID AS StateId, J.STATE_NAME AS StateName, J.INVOCATION_DATA AS InvocationData, + J.ARGUMENTS AS Arguments, J.CREATED_AT AS CreatedAt, J.EXPIRE_AT AS ExpireAt, + S.REASON AS StateReason, S.DATA AS StateData + FROM {T("Job")} J + LEFT JOIN {T("JobState")} S ON S.ID = J.STATE_ID"; } var jobs = connection.Query(enqueuedJobsSql, new { JOB_IDS = enumerable }).ToList(); @@ -482,13 +482,15 @@ LEFT JOIN HF_JOB_STATE S }); } - private static JobList FetchedJobs(IDbConnection connection, IEnumerable jobIds) + private JobList FetchedJobs(IDbConnection connection, IEnumerable jobIds) { - const string fetchedJobsSql = @" - SELECT J.ID AS Id, J.STATE_ID AS StateId, J.STATE_NAME AS StateName, J.INVOCATION_DATA AS InvocationData, J.ARGUMENTS AS Arguments, J.CREATED_AT AS CreatedAt, J.EXPIRE_AT AS ExpireAt, S.REASON AS StateReason, S.DATA AS StateData - FROM HF_JOB J - LEFT JOIN HF_JOB_STATE S ON S.ID = J.STATE_ID - WHERE J.ID IN :JOB_IDS"; + var fetchedJobsSql = $@" + SELECT J.ID AS Id, J.STATE_ID AS StateId, J.STATE_NAME AS StateName, J.INVOCATION_DATA AS InvocationData, + J.ARGUMENTS AS Arguments, J.CREATED_AT AS CreatedAt, J.EXPIRE_AT AS ExpireAt, + S.REASON AS StateReason, S.DATA AS StateData + FROM {T("Job")} J + LEFT JOIN {T("JobState")} S ON S.ID = J.STATE_ID + WHERE J.ID IN :JOB_IDS"; var jobs = connection.Query(fetchedJobsSql, new { JOB_IDS = jobIds }).ToList(); @@ -506,7 +508,7 @@ FROM HF_JOB J return new JobList(result); } - private static Dictionary GetHourlyTimelineStats(IDbConnection connection, string type) + private Dictionary GetHourlyTimelineStats(IDbConnection connection, string type) { var endDate = DateTime.UtcNow; var dates = new List(); diff --git a/Hangfire.Oracle/OracleDistributedLock.cs b/Hangfire.Oracle/OracleDistributedLock.cs index 69bf0d2..93d322a 100644 --- a/Hangfire.Oracle/OracleDistributedLock.cs +++ b/Hangfire.Oracle/OracleDistributedLock.cs @@ -1,9 +1,7 @@ using System; using System.Data; using System.Threading; - using Dapper; - using Hangfire.Logging; namespace Hangfire.Oracle.Core @@ -15,17 +13,24 @@ public class OracleDistributedLock : IDisposable, IComparable private readonly OracleStorage _storage; private readonly DateTime _start; private readonly CancellationToken _cancellationToken; + private readonly IDbConnection _connection; + private readonly bool _ownsConnection; private const int DelayBetweenPasses = 100; public OracleDistributedLock(OracleStorage storage, string resource, TimeSpan timeout) - : this(storage.CreateAndOpenConnection(), resource, timeout) { - _storage = storage; + Logger.TraceFormat("OracleDistributedLock resource={0}, timeout={1}", resource, timeout); + + _storage = storage ?? throw new ArgumentNullException(nameof(storage)); + _connection = storage.CreateAndOpenConnection(); + Resource = resource; + _timeout = timeout; + _cancellationToken = default; + _start = DateTime.UtcNow; + _ownsConnection = true; } - private readonly IDbConnection _connection; - public OracleDistributedLock(IDbConnection connection, string resource, TimeSpan timeout) : this(connection, resource, timeout, new CancellationToken()) { @@ -35,42 +40,73 @@ public OracleDistributedLock(IDbConnection connection, string resource, TimeSpan { Logger.TraceFormat("OracleDistributedLock resource={0}, timeout={1}", resource, timeout); + if (connection == null) throw new ArgumentNullException(nameof(connection)); + + _storage = ExtractStorageFromConnection(connection); + _connection = connection; Resource = resource; _timeout = timeout; - _connection = connection; _cancellationToken = cancellationToken; _start = DateTime.UtcNow; + _ownsConnection = false; } public string Resource { get; } + private OracleStorage ExtractStorageFromConnection(IDbConnection connection) + { + var storage = OracleStorageConnectionContext.Current; + if (storage == null) + { + Logger.Warn("Cannot extract OracleStorage from connection context. Using default table names."); + } + return storage; + } + + private string GetDistributedLockTable() + { + var tableName = _storage?.TableNameProvider?.GetTableName("DistributedLock") ?? "HF_DISTRIBUTED_LOCK"; + Logger.DebugFormat("GetDistributedLockTable resolved to: {0}", tableName); + return tableName; + } + private int AcquireLock(string resource, TimeSpan timeout) { - return - _connection - .Execute( - @" -INSERT INTO HF_DISTRIBUTED_LOCK (""RESOURCE"", CREATED_AT) - (SELECT :RES, :NOW - FROM DUAL - WHERE NOT EXISTS - (SELECT ""RESOURCE"", CREATED_AT - FROM HF_DISTRIBUTED_LOCK - WHERE ""RESOURCE"" = :RES AND CREATED_AT > :EXPIRED)) -", - new - { - RES = resource, - NOW = DateTime.UtcNow, - EXPIRED = DateTime.UtcNow.Add(timeout.Negate()) - }); + var tableName = GetDistributedLockTable(); + + if (string.IsNullOrWhiteSpace(tableName)) + { + throw new InvalidOperationException("DistributedLock table name is empty!"); + } + + var sql = $@"INSERT INTO {tableName} (""RESOURCE"", CREATED_AT) + SELECT :RES, :NOW + FROM DUAL + WHERE NOT EXISTS ( + SELECT 1 FROM {tableName} + WHERE ""RESOURCE"" = :RES + AND CREATED_AT > :EXPIRED + )"; + + Logger.DebugFormat("AcquireLock SQL:\n{0}", sql); + + return _connection.Execute( + sql, + new { + RES = resource, + NOW = DateTime.UtcNow, + EXPIRED = DateTime.UtcNow.Add(timeout.Negate()) + }); } public void Dispose() { Release(); - _storage?.ReleaseConnection(_connection); + if (_ownsConnection) + { + _storage?.ReleaseConnection(_connection); + } } internal OracleDistributedLock Acquire() @@ -107,16 +143,18 @@ internal void Release() { Logger.TraceFormat("Release resource={0}", Resource); - _connection - .Execute( - @" -DELETE FROM HF_DISTRIBUTED_LOCK - WHERE ""RESOURCE"" = :RES -", - new - { - RES = Resource - }); + var tableName = GetDistributedLockTable(); + if (string.IsNullOrWhiteSpace(tableName)) + { + Logger.Error("DistributedLock table name is empty during Release!"); + return; + } + + var sql = $@"DELETE FROM {tableName} WHERE ""RESOURCE"" = :RES"; + + Logger.DebugFormat("Release SQL:\n{0}", sql); + + _connection.Execute(sql, new { RES = Resource }); } public int CompareTo(object obj) @@ -130,8 +168,21 @@ public int CompareTo(object obj) { return string.Compare(Resource, oracleDistributedLock.Resource, StringComparison.OrdinalIgnoreCase); } - + throw new ArgumentException("Object is not a OracleDistributedLock"); } } + + // Thread-local context to pass storage through connection calls + internal static class OracleStorageConnectionContext + { + [ThreadStatic] + private static OracleStorage _current; + + public static OracleStorage Current + { + get => _current; + set => _current = value; + } + } } \ No newline at end of file diff --git a/Hangfire.Oracle/OracleObjectsInstaller.cs b/Hangfire.Oracle/OracleObjectsInstaller.cs index 23046f5..58267d7 100644 --- a/Hangfire.Oracle/OracleObjectsInstaller.cs +++ b/Hangfire.Oracle/OracleObjectsInstaller.cs @@ -1,23 +1,24 @@ using System; using System.Data; -using System.IO; -using System.Linq; -using System.Reflection; - +using System.Text; using Dapper; - using Hangfire.Logging; +using Hangfire.Oracle.Core.Configuration; namespace Hangfire.Oracle.Core { public static class OracleObjectsInstaller { private static readonly ILog Log = LogProvider.GetLogger(typeof(OracleStorage)); - public static void Install(IDbConnection connection, string schemaName) + + public static void Install(IDbConnection connection, HangfireTableNameProvider tableNameProvider) { if (connection == null) throw new ArgumentNullException(nameof(connection)); + if (tableNameProvider == null) throw new ArgumentNullException(nameof(tableNameProvider)); - if (TablesExists(connection, schemaName)) + var schema = tableNameProvider.GetSchemaName(); + + if (TablesExists(connection, schema, tableNameProvider)) { Log.Info("DB tables already exist. Exit install"); return; @@ -25,56 +26,268 @@ public static void Install(IDbConnection connection, string schemaName) Log.Info("Start installing Hangfire SQL objects..."); - var script = GetStringResource("Hangfire.Oracle.Core.Install.sql"); - + var script = GenerateInstallScript(tableNameProvider); var sqlCommands = script.Split(new[] { ';' }, StringSplitOptions.RemoveEmptyEntries); - sqlCommands.ToList().ForEach(s => connection.Execute(s)); - Log.Info("Hangfire SQL objects installed."); + foreach (var sqlCommand in sqlCommands) + { + var trimmedCommand = sqlCommand.Trim(); + if (!string.IsNullOrWhiteSpace(trimmedCommand)) + { + connection.Execute(trimmedCommand); + } + } + + Log.Info("Hangfire SQL objects installed."); } - private static bool TablesExists(IDbConnection connection, string schemaName) + private static bool TablesExists(IDbConnection connection, string schemaName, HangfireTableNameProvider tableNameProvider) { string tableExistsQuery; + var jobTableName = tableNameProvider.GetTableName("Job"); if (!string.IsNullOrEmpty(schemaName)) { - tableExistsQuery = $@"SELECT TABLE_NAME -FROM all_tables -WHERE OWNER = '{schemaName}' AND TABLE_NAME LIKE 'HF_%' -ORDER BY TABLE_NAME"; + tableExistsQuery = $@" + SELECT TABLE_NAME + FROM all_tables + WHERE OWNER = '{schemaName}' AND TABLE_NAME = '{jobTableName}'"; } else { - tableExistsQuery = @"SELECT TABLE_NAME -FROM all_tables -WHERE TABLE_NAME LIKE 'HF_%' -ORDER BY TABLE_NAME"; + tableExistsQuery = $@" + SELECT TABLE_NAME + FROM all_tables + WHERE TABLE_NAME = '{jobTableName}'"; } return connection.ExecuteScalar(tableExistsQuery) != null; } - private static string GetStringResource(string resourceName) + private static string GenerateInstallScript(HangfireTableNameProvider tableNameProvider) { -#if NET45 - var assembly = typeof(OracleObjectsInstaller).Assembly; -#else - var assembly = typeof(OracleObjectsInstaller).GetTypeInfo().Assembly; -#endif + var sb = new StringBuilder(); - using (var stream = assembly.GetManifestResourceStream(resourceName)) - { - if (stream == null) - { - throw new InvalidOperationException($"Requested resource `{resourceName}` was not found in the assembly `{assembly}`."); - } + var primarySequence = tableNameProvider.GetPrimarySequenceName(); + var jobIdSequence = tableNameProvider.GetJobIdSequenceName(); + var prefix = tableNameProvider.GetInstancePrefix(); - using (var reader = new StreamReader(stream)) - { - return reader.ReadToEnd(); - } - } + var jobTable = tableNameProvider.GetTableName("Job"); + var counterTable = tableNameProvider.GetTableName("Counter"); + var aggregatedCounterTable = tableNameProvider.GetTableName("AggregatedCounter"); + var distributedLockTable = tableNameProvider.GetTableName("DistributedLock"); + var hashTable = tableNameProvider.GetTableName("Hash"); + var jobParameterTable = tableNameProvider.GetTableName("JobParameter"); + var jobQueueTable = tableNameProvider.GetTableName("JobQueue"); + var jobStateTable = tableNameProvider.GetTableName("JobState"); + var serverTable = tableNameProvider.GetTableName("Server"); + var setTable = tableNameProvider.GetTableName("Set"); + var listTable = tableNameProvider.GetTableName("List"); + + sb.AppendLine($"CREATE SEQUENCE {primarySequence} START WITH 1 MAXVALUE 9999999999999999999999999999 MINVALUE 1 NOCYCLE CACHE 20 NOORDER;"); + sb.AppendLine(); + sb.AppendLine($"CREATE SEQUENCE {jobIdSequence} START WITH 1 MAXVALUE 9999999999999999999999999999 MINVALUE 1 NOCYCLE CACHE 20 NOORDER;"); + sb.AppendLine(); + + sb.AppendLine($@"CREATE TABLE {jobTable} ( + ID NUMBER(10), + STATE_ID NUMBER(10), + STATE_NAME NVARCHAR2(20), + INVOCATION_DATA NCLOB, + ARGUMENTS NCLOB, + CREATED_AT TIMESTAMP(4), + EXPIRE_AT TIMESTAMP(4) + ) + LOB (INVOCATION_DATA) STORE AS BASICFILE + (ENABLE STORAGE IN ROW CHUNK 8192 RETENTION NOCACHE LOGGING) + LOB (ARGUMENTS) STORE AS BASICFILE + (ENABLE STORAGE IN ROW CHUNK 8192 RETENTION NOCACHE LOGGING) + LOGGING NOCOMPRESS NOCACHE NOPARALLEL MONITORING;"); + sb.AppendLine(); + + sb.AppendLine($@"ALTER TABLE {jobTable} ADD ( + CONSTRAINT {tableNameProvider.GetConstraintName(jobTable, "PK")} PRIMARY KEY (ID) + USING INDEX (CREATE UNIQUE INDEX {tableNameProvider.GetIndexName(jobTable, "PK")} ON {jobTable} (ID)) + ENABLE VALIDATE);"); + sb.AppendLine(); + + sb.AppendLine($@"CREATE TABLE {counterTable} ( + ID NUMBER(10), + KEY NVARCHAR2(255), + VALUE NUMBER(10), + EXPIRE_AT TIMESTAMP(4) + ) + LOGGING NOCOMPRESS NOCACHE NOPARALLEL MONITORING;"); + sb.AppendLine(); + + sb.AppendLine($@"ALTER TABLE {counterTable} ADD ( + CONSTRAINT {tableNameProvider.GetConstraintName(counterTable, "PK")} PRIMARY KEY (ID) + USING INDEX (CREATE UNIQUE INDEX {tableNameProvider.GetIndexName(counterTable, "PK")} ON {counterTable} (ID)) + ENABLE VALIDATE);"); + sb.AppendLine(); + + sb.AppendLine($@"CREATE TABLE {aggregatedCounterTable} ( + ID NUMBER(10), + KEY NVARCHAR2(255), + VALUE NUMBER(10), + EXPIRE_AT TIMESTAMP(4) + ) + LOGGING NOCOMPRESS NOCACHE NOPARALLEL MONITORING;"); + sb.AppendLine(); + + sb.AppendLine($@"ALTER TABLE {aggregatedCounterTable} ADD ( + CONSTRAINT {tableNameProvider.GetConstraintName(aggregatedCounterTable, "PK")} PRIMARY KEY (ID) + USING INDEX (CREATE UNIQUE INDEX {tableNameProvider.GetIndexName(aggregatedCounterTable, "PK")} ON {aggregatedCounterTable} (ID)) + ENABLE VALIDATE, + CONSTRAINT {tableNameProvider.GetConstraintName(aggregatedCounterTable, "UQ")} UNIQUE (KEY) + USING INDEX (CREATE UNIQUE INDEX {tableNameProvider.GetIndexName(aggregatedCounterTable, "UQ")} ON {aggregatedCounterTable} (KEY)) + ENABLE VALIDATE);"); + sb.AppendLine(); + + sb.AppendLine($@"CREATE TABLE {distributedLockTable} ( + ""RESOURCE"" NVARCHAR2(100), + CREATED_AT TIMESTAMP(4)) + LOGGING NOCOMPRESS NOCACHE NOPARALLEL MONITORING;"); + sb.AppendLine(); + + sb.AppendLine($@"CREATE TABLE {hashTable} ( + ID NUMBER(10), + KEY NVARCHAR2(255), + VALUE NCLOB, + EXPIRE_AT TIMESTAMP(4), + FIELD NVARCHAR2(40) + ) + LOB (VALUE) STORE AS BASICFILE + (ENABLE STORAGE IN ROW CHUNK 8192 RETENTION NOCACHE LOGGING) + LOGGING NOCOMPRESS NOCACHE NOPARALLEL MONITORING;"); + sb.AppendLine(); + + sb.AppendLine($@"ALTER TABLE {hashTable} ADD ( + CONSTRAINT {tableNameProvider.GetConstraintName(hashTable, "PK")} PRIMARY KEY (ID) + USING INDEX (CREATE UNIQUE INDEX {tableNameProvider.GetIndexName(hashTable, "PK")} ON {hashTable} (ID)) + ENABLE VALIDATE, + CONSTRAINT {tableNameProvider.GetConstraintName(hashTable, "UQ")} UNIQUE (KEY, FIELD) + USING INDEX (CREATE UNIQUE INDEX {tableNameProvider.GetIndexName(hashTable, "UQ")} ON {hashTable} (KEY, FIELD)) + ENABLE VALIDATE);"); + sb.AppendLine(); + + sb.AppendLine($@"CREATE TABLE {jobParameterTable} ( + ID NUMBER(10), + NAME NVARCHAR2(40), + VALUE NCLOB, + JOB_ID NUMBER(10) + ) + LOB (VALUE) STORE AS BASICFILE + (ENABLE STORAGE IN ROW CHUNK 8192 RETENTION NOCACHE LOGGING) + LOGGING NOCOMPRESS NOCACHE NOPARALLEL MONITORING;"); + sb.AppendLine(); + + sb.AppendLine($@"ALTER TABLE {jobParameterTable} ADD ( + CONSTRAINT {tableNameProvider.GetConstraintName(jobParameterTable, "PK")} PRIMARY KEY (ID) + USING INDEX (CREATE UNIQUE INDEX {tableNameProvider.GetIndexName(jobParameterTable, "PK")} ON {jobParameterTable} (ID)) + ENABLE VALIDATE);"); + sb.AppendLine(); + + sb.AppendLine($@"ALTER TABLE {jobParameterTable} ADD ( + CONSTRAINT {tableNameProvider.GetConstraintName(jobParameterTable, "FK_JOB")} FOREIGN KEY (JOB_ID) + REFERENCES {jobTable} (ID) + ON DELETE CASCADE ENABLE VALIDATE);"); + sb.AppendLine(); + + sb.AppendLine($@"CREATE TABLE {jobQueueTable} ( + ID NUMBER(10), + JOB_ID NUMBER(10), + QUEUE NVARCHAR2(50), + FETCHED_AT TIMESTAMP(4), + FETCH_TOKEN NVARCHAR2(36)) + LOGGING NOCOMPRESS NOCACHE NOPARALLEL MONITORING;"); + sb.AppendLine(); + + sb.AppendLine($@"ALTER TABLE {jobQueueTable} ADD ( + CONSTRAINT {tableNameProvider.GetConstraintName(jobQueueTable, "PK")} PRIMARY KEY (ID) + USING INDEX (CREATE UNIQUE INDEX {tableNameProvider.GetIndexName(jobQueueTable, "PK")} ON {jobQueueTable} (ID)) + ENABLE VALIDATE);"); + sb.AppendLine(); + + sb.AppendLine($@"ALTER TABLE {jobQueueTable} ADD ( + CONSTRAINT {tableNameProvider.GetConstraintName(jobQueueTable, "FK_JOB")} FOREIGN KEY (JOB_ID) + REFERENCES {jobTable} (ID) + ON DELETE CASCADE ENABLE VALIDATE);"); + sb.AppendLine(); + + sb.AppendLine($@"CREATE TABLE {jobStateTable} ( + ID NUMBER(10), + JOB_ID NUMBER(10), + NAME NVARCHAR2(20), + REASON NVARCHAR2(100), + CREATED_AT TIMESTAMP(4), + DATA NCLOB) + LOB (DATA) STORE AS BASICFILE + (ENABLE STORAGE IN ROW CHUNK 8192 RETENTION NOCACHE LOGGING) + LOGGING NOCOMPRESS NOCACHE NOPARALLEL MONITORING;"); + sb.AppendLine(); + + sb.AppendLine($@"ALTER TABLE {jobStateTable} ADD ( + CONSTRAINT {tableNameProvider.GetConstraintName(jobStateTable, "PK")} PRIMARY KEY (ID) + USING INDEX (CREATE UNIQUE INDEX {tableNameProvider.GetIndexName(jobStateTable, "PK")} ON {jobStateTable} (ID)) + ENABLE VALIDATE);"); + sb.AppendLine(); + + sb.AppendLine($@"ALTER TABLE {jobStateTable} ADD ( + CONSTRAINT {tableNameProvider.GetConstraintName(jobStateTable, "FK_JOB")} FOREIGN KEY (JOB_ID) + REFERENCES {jobTable} (ID) + ON DELETE CASCADE ENABLE VALIDATE);"); + sb.AppendLine(); + + sb.AppendLine($@"CREATE TABLE {serverTable} ( + ID NVARCHAR2(100), + DATA NCLOB, + LAST_HEART_BEAT TIMESTAMP(4)) + LOB (DATA) STORE AS BASICFILE + (ENABLE STORAGE IN ROW CHUNK 8192 RETENTION NOCACHE LOGGING) + LOGGING NOCOMPRESS NOCACHE NOPARALLEL MONITORING;"); + sb.AppendLine(); + + sb.AppendLine($@"ALTER TABLE {serverTable} ADD ( + CONSTRAINT {tableNameProvider.GetConstraintName(serverTable, "PK")} PRIMARY KEY (ID) + USING INDEX (CREATE UNIQUE INDEX {tableNameProvider.GetIndexName(serverTable, "PK")} ON {serverTable} (ID)) + ENABLE VALIDATE);"); + sb.AppendLine(); + + sb.AppendLine($@"CREATE TABLE {setTable} ( + ID NUMBER(10), + KEY NVARCHAR2(255), + VALUE NVARCHAR2(255), + SCORE FLOAT(126), + EXPIRE_AT TIMESTAMP(4)) + LOGGING NOCOMPRESS NOCACHE NOPARALLEL MONITORING;"); + sb.AppendLine(); + + sb.AppendLine($@"ALTER TABLE {setTable} ADD ( + CONSTRAINT {tableNameProvider.GetConstraintName(setTable, "PK")} PRIMARY KEY (ID) + USING INDEX (CREATE UNIQUE INDEX {tableNameProvider.GetIndexName(setTable, "PK")} ON {setTable} (ID)) + ENABLE VALIDATE, + CONSTRAINT {tableNameProvider.GetConstraintName(setTable, "UQ")} UNIQUE (KEY, VALUE) + USING INDEX (CREATE UNIQUE INDEX {tableNameProvider.GetIndexName(setTable, "UQ")} ON {setTable} (KEY, VALUE)) + ENABLE VALIDATE);"); + sb.AppendLine(); + + sb.AppendLine($@"CREATE TABLE {listTable} ( + ID NUMBER(10), + KEY NVARCHAR2(255), + VALUE NCLOB, + EXPIRE_AT TIMESTAMP(4)) + LOB (VALUE) STORE AS BASICFILE + (ENABLE STORAGE IN ROW CHUNK 8192 RETENTION NOCACHE LOGGING) + LOGGING NOCOMPRESS NOCACHE NOPARALLEL MONITORING;"); + sb.AppendLine(); + + sb.AppendLine($@"ALTER TABLE {listTable} ADD ( + CONSTRAINT {tableNameProvider.GetConstraintName(listTable, "PK")} PRIMARY KEY (ID) + USING INDEX (CREATE UNIQUE INDEX {tableNameProvider.GetIndexName(listTable, "PK")} ON {listTable} (ID)) + ENABLE VALIDATE);"); + + return sb.ToString(); } } } diff --git a/Hangfire.Oracle/OracleStorage.cs b/Hangfire.Oracle/OracleStorage.cs index cfc91af..07c2bae 100644 --- a/Hangfire.Oracle/OracleStorage.cs +++ b/Hangfire.Oracle/OracleStorage.cs @@ -3,17 +3,16 @@ using System.Data; using System.Linq; using System.Text; - using Dapper; - +using Hangfire; using Hangfire.Annotations; using Hangfire.Logging; -using Hangfire.Oracle.Core.JobQueue; -using Hangfire.Oracle.Core.Monitoring; using Hangfire.Server; using Hangfire.Storage; - using Oracle.ManagedDataAccess.Client; +using Hangfire.Oracle.Core.Configuration; +using Hangfire.Oracle.Core.JobQueue; +using Hangfire.Oracle.Core.Monitoring; namespace Hangfire.Oracle.Core { @@ -27,6 +26,8 @@ public class OracleStorage : JobStorage, IDisposable private readonly OracleStorageOptions _options; public virtual PersistentJobQueueProviderCollection QueueProviders { get; private set; } + + public HangfireTableNameProvider TableNameProvider { get; private set; } public OracleStorage(string connectionString) : this(connectionString, new OracleStorageOptions()) @@ -50,6 +51,8 @@ public OracleStorage(string connectionString, OracleStorageOptions options) } _options = options ?? throw new ArgumentNullException(nameof(options)); + + InitializeTableNameProvider(); PrepareSchemaIfNecessary(options); InitializeQueueProviders(); @@ -60,18 +63,30 @@ public OracleStorage(Func connectionFactory, OracleStorageOptions _connectionFactory = connectionFactory ?? throw new ArgumentNullException(nameof(connectionFactory)); _options = options ?? throw new ArgumentNullException(nameof(options)); + + InitializeTableNameProvider(); PrepareSchemaIfNecessary(options); InitializeQueueProviders(); } + private void InitializeTableNameProvider() + { + var config = _options.InstanceConfiguration ?? new HangfireConfiguration(); + + config.Tables = config.Tables ?? new Dictionary(); + config.Sequence = config.Sequence ?? new SequenceConfiguration(); + + TableNameProvider = new HangfireTableNameProvider(config); + } + private void PrepareSchemaIfNecessary(OracleStorageOptions options) { if (options.PrepareSchemaIfNecessary) { using (var connection = CreateAndOpenConnection()) { - OracleObjectsInstaller.Install(connection, options.SchemaName); + OracleObjectsInstaller.Install(connection, TableNameProvider); } } } @@ -91,7 +106,7 @@ public override IEnumerable GetComponents() public override void WriteOptionsToLog(ILog logger) { - logger.Info("Using the following options for SQL Server job storage:"); + logger.Info("Using the following options for Oracle job storage:"); logger.InfoFormat(" Queue poll interval: {0}.", _options.QueuePollInterval); } @@ -170,7 +185,6 @@ internal void UseTransaction([InstantHandle] Action action) return true; }, null); } - internal T UseTransaction([InstantHandle] Func func, IsolationLevel? isolationLevel) { return UseConnection(connection => @@ -184,7 +198,6 @@ internal T UseTransaction([InstantHandle] Func func, Isolat } }); } - internal void UseConnection([InstantHandle] Action action) { UseConnection(connection => @@ -200,11 +213,14 @@ internal T UseConnection([InstantHandle] Func func) try { + OracleStorageConnectionContext.Current = this; + connection = CreateAndOpenConnection(); return func(connection); } finally { + OracleStorageConnectionContext.Current = null; ReleaseConnection(connection); } } @@ -217,19 +233,20 @@ internal IDbConnection CreateAndOpenConnection() { connection.Open(); - if (!string.IsNullOrWhiteSpace(_options.SchemaName)) + var schema = TableNameProvider.GetSchemaName(); + if (!string.IsNullOrWhiteSpace(schema)) { - connection.Execute($"ALTER SESSION SET CURRENT_SCHEMA={_options.SchemaName}"); + connection.Execute($"ALTER SESSION SET CURRENT_SCHEMA={schema}"); } } return connection; } - internal void ReleaseConnection(IDbConnection connection) { connection?.Dispose(); } + public void Dispose() { } diff --git a/Hangfire.Oracle/OracleStorageConnection.cs b/Hangfire.Oracle/OracleStorageConnection.cs index c0af426..f6b145e 100644 --- a/Hangfire.Oracle/OracleStorageConnection.cs +++ b/Hangfire.Oracle/OracleStorageConnection.cs @@ -3,10 +3,8 @@ using System.Data; using System.Linq; using System.Threading; - using Dapper; using Dapper.Oracle; - using Hangfire.Common; using Hangfire.Logging; using Hangfire.Oracle.Core.Entities; @@ -20,11 +18,16 @@ public class OracleStorageConnection : JobStorageConnection private static readonly ILog Logger = LogProvider.GetLogger(typeof(OracleStorageConnection)); private readonly OracleStorage _storage; + public OracleStorageConnection(OracleStorage storage) { _storage = storage ?? throw new ArgumentNullException(nameof(storage)); } + private string T(string logicalName) => _storage.TableNameProvider.GetTableName(logicalName); + private string GetPrimarySequence() => _storage.TableNameProvider.GetPrimarySequenceName(); + private string GetJobIdSequence() => _storage.TableNameProvider.GetJobIdSequenceName(); + public override IWriteOnlyTransaction CreateWriteTransaction() { return new OracleWriteOnlyTransaction(_storage); @@ -55,7 +58,7 @@ public override string CreateExpiredJob(Job job, IDictionary par return _storage.UseConnection(connection => { - var jobId = connection.GetNextJobId(); + var jobId = connection.GetNextJobId(GetJobIdSequence()); var oracleDynamicParameters = new OracleDynamicParameters(); oracleDynamicParameters.AddDynamicParams(new @@ -65,13 +68,11 @@ public override string CreateExpiredJob(Job job, IDictionary par EXPIRE_AT = createdAt.Add(expireIn) }); oracleDynamicParameters.Add("INVOCATION_DATA", SerializationHelper.Serialize(invocationData, SerializationOption.User), OracleMappingType.NClob, ParameterDirection.Input); - oracleDynamicParameters.Add("ARGUMENTS", arguments.Arguments, OracleMappingType.NClob, ParameterDirection.Input); + oracleDynamicParameters.Add("ARGUMENTS", arguments.Arguments, OracleMappingType.NClob , ParameterDirection.Input); connection.Execute( - @" - INSERT INTO HF_JOB (ID, INVOCATION_DATA, ARGUMENTS, CREATED_AT, EXPIRE_AT) - VALUES (:ID, :INVOCATION_DATA, :ARGUMENTS, :CREATED_AT, :EXPIRE_AT) -", + $@"INSERT INTO {T("Job")} (ID, INVOCATION_DATA, ARGUMENTS, CREATED_AT, EXPIRE_AT) + VALUES (:ID, :INVOCATION_DATA, :ARGUMENTS, :CREATED_AT, :EXPIRE_AT)", oracleDynamicParameters); if (parameters.Count > 0) @@ -86,12 +87,12 @@ INSERT INTO HF_JOB (ID, INVOCATION_DATA, ARGUMENTS, CREATED_AT, EXPIRE_AT) JOB_ID = jobId, NAME = parameter.Key }); - dynamicParameters.Add("VALUE", parameter.Value, OracleMappingType.NClob, ParameterDirection.Input); + dynamicParameters.Add("VALUE", parameter.Value,OracleMappingType.NClob, ParameterDirection.Input); parameterArray[parameterIndex++] = dynamicParameters; } - connection.Execute(@"INSERT INTO HF_JOB_PARAMETER (ID, NAME, VALUE, JOB_ID) VALUES (HF_SEQUENCE.NEXTVAL, :NAME, :VALUE, :JOB_ID)", parameterArray); + connection.Execute($"INSERT INTO {T("JobParameter")} (ID, NAME, VALUE, JOB_ID) VALUES ({GetPrimarySequence()}.NEXTVAL, :NAME, :VALUE, :JOB_ID)", parameterArray); } return jobId.ToString(); @@ -136,18 +137,16 @@ public override void SetJobParameter(string id, string name, string value) { var oracleDynamicParameters = new OracleDynamicParameters(); oracleDynamicParameters.AddDynamicParams(new { JOB_ID = id, NAME = name }); - oracleDynamicParameters.Add("VALUE", value, OracleMappingType.NClob, ParameterDirection.Input); + oracleDynamicParameters.Add("VALUE", value,OracleMappingType.NClob, ParameterDirection.Input); connection.Execute( - @" - MERGE INTO HF_JOB_PARAMETER JP - USING (SELECT 1 FROM DUAL) SRC - ON (JP.NAME = :NAME AND JP.JOB_ID = :JOB_ID) - WHEN MATCHED THEN - UPDATE SET VALUE = :VALUE - WHEN NOT MATCHED THEN - INSERT (ID, JOB_ID, NAME, VALUE) - VALUES (HF_SEQUENCE.NEXTVAL, :JOB_ID, :NAME, :VALUE) -", + $@"MERGE INTO {T("JobParameter")} JP + USING (SELECT 1 FROM DUAL) SRC + ON (JP.NAME = :NAME AND JP.JOB_ID = :JOB_ID) + WHEN MATCHED THEN + UPDATE SET VALUE = :VALUE + WHEN NOT MATCHED THEN + INSERT (ID, JOB_ID, NAME, VALUE) + VALUES ({GetPrimarySequence()}.NEXTVAL, :JOB_ID, :NAME, :VALUE)", oracleDynamicParameters); }); } @@ -166,9 +165,9 @@ public override string GetJobParameter(string id, string name) return _storage.UseConnection(connection => connection.QuerySingleOrDefault( - "SELECT VALUE as Value " + - " FROM HF_JOB_PARAMETER " + - " WHERE JOB_ID = :ID AND NAME = :NAME", + $@"SELECT VALUE as Value + FROM {T("JobParameter")} + WHERE JOB_ID = :ID AND NAME = :NAME", new { ID = id, NAME = name })); } @@ -182,9 +181,9 @@ public override JobData GetJobData(string jobId) return _storage.UseConnection(connection => { var jobData = connection.QuerySingleOrDefault( - "SELECT INVOCATION_DATA AS InvocationData, STATE_NAME AS StateName, ARGUMENTS AS Arguments, CREATED_AT AS CreatedAt " + - " FROM HF_JOB " + - " WHERE ID = :ID", + $@"SELECT INVOCATION_DATA AS InvocationData, STATE_NAME AS StateName, ARGUMENTS AS Arguments, CREATED_AT AS CreatedAt + FROM {T("Job")} + WHERE ID = :ID", new { ID = jobId }); if (jobData == null) @@ -227,11 +226,10 @@ public override StateData GetStateData(string jobId) return _storage.UseConnection(connection => { var sqlState = connection.QuerySingleOrDefault( - " SELECT S.NAME AS Name, S.REASON AS Reason, S.DATA AS Data " + - " FROM HF_JOB_STATE S" + - " INNER JOIN HF_JOB J" + - " ON J.STATE_ID = S.ID " + - " WHERE J.ID = :JOB_ID", + $@"SELECT S.NAME AS Name, S.REASON AS Reason, S.DATA AS Data + FROM {T("JobState")} S + INNER JOIN {T("Job")} J ON J.STATE_ID = S.ID + WHERE J.ID = :JOB_ID", new { JOB_ID = jobId }); if (sqlState == null) @@ -267,16 +265,14 @@ public override void AnnounceServer(string serverId, ServerContext context) _storage.UseConnection(connection => { connection.Execute( - @" - MERGE INTO HF_SERVER S - USING (SELECT 1 FROM DUAL) SRC - ON (S.ID = :ID) - WHEN MATCHED THEN - UPDATE SET LAST_HEART_BEAT = :LAST_HEART_BEAT - WHEN NOT MATCHED THEN - INSERT (ID, DATA, LAST_HEART_BEAT) - VALUES (:ID, :DATA, :LAST_HEART_BEAT) -", + $@"MERGE INTO {T("Server")} S + USING (SELECT 1 FROM DUAL) SRC + ON (S.ID = :ID) + WHEN MATCHED THEN + UPDATE SET LAST_HEART_BEAT = :LAST_HEART_BEAT + WHEN NOT MATCHED THEN + INSERT (ID, DATA, LAST_HEART_BEAT) + VALUES (:ID, :DATA, :LAST_HEART_BEAT)", new { ID = serverId, @@ -301,7 +297,7 @@ public override void RemoveServer(string serverId) _storage.UseConnection(connection => { connection.Execute( - "DELETE FROM HF_SERVER where ID = :ID", + $"DELETE FROM {T("Server")} where ID = :ID", new { ID = serverId }); }); } @@ -316,9 +312,9 @@ public override void Heartbeat(string serverId) _storage.UseConnection(connection => { connection.Execute( - " UPDATE HF_SERVER" + - " SET LAST_HEART_BEAT = :NOW" + - " WHERE ID = :ID", + $@"UPDATE {T("Server")} + SET LAST_HEART_BEAT = :NOW + WHERE ID = :ID", new { NOW = DateTime.UtcNow, ID = serverId }); }); } @@ -333,8 +329,8 @@ public override int RemoveTimedOutServers(TimeSpan timeOut) return _storage.UseConnection(connection => connection.Execute( - " DELETE FROM HF_SERVER" + - " WHERE LAST_HEART_BEAT < :TIME_OUT_AT", + $@"DELETE FROM {T("Server")} + WHERE LAST_HEART_BEAT < :TIME_OUT_AT", new { TIME_OUT_AT = DateTime.UtcNow.Add(timeOut.Negate()) })); } @@ -345,9 +341,9 @@ public override long GetSetCount(string key) return _storage.UseConnection(connection => connection.QueryFirst( - "SELECT COUNT(KEY) " + - " FROM HF_SET" + - " WHERE KEY = :KEY", + $@"SELECT COUNT(KEY) + FROM {T("Set")} + WHERE KEY = :KEY", new { KEY = key })); } @@ -359,13 +355,12 @@ public override List GetRangeFromSet(string key, int startingFrom, int e } return _storage.UseConnection(connection => - connection.Query(@" -SELECT VALUE as Value - FROM (SELECT VALUE, RANK () OVER (ORDER BY ID) AS RANK - FROM HF_SET - WHERE KEY = :KEY) - WHERE RANK BETWEEN :S AND :E -", + connection.Query($@" + SELECT VALUE as Value + FROM (SELECT VALUE, RANK () OVER (ORDER BY ID) AS RANK + FROM {T("Set")} + WHERE KEY = :KEY) + WHERE RANK BETWEEN :S AND :E", new { KEY = key, S = startingFrom + 1, E = endingAt + 1 }).ToList()); } @@ -380,9 +375,9 @@ public override HashSet GetAllItemsFromSet(string key) _storage.UseConnection(connection => { var result = connection.Query( - "SELECT VALUE AS Value" + - " FROM HF_SET" + - " WHERE KEY = :KEY", + $@"SELECT VALUE AS Value + FROM {T("Set")} + WHERE KEY = :KEY", new { KEY = key }); return new HashSet(result); @@ -404,14 +399,12 @@ public override string GetFirstByLowestScoreFromSet(string key, double fromScore return _storage.UseConnection(connection => connection.QuerySingleOrDefault( - @" -SELECT * - FROM ( SELECT VALUE AS Value - FROM HF_SET - WHERE KEY = :KEY AND SCORE BETWEEN :F AND :T - ORDER BY SCORE) - WHERE ROWNUM = 1 -", + $@"SELECT * + FROM (SELECT VALUE AS Value + FROM {T("Set")} + WHERE KEY = :KEY AND SCORE BETWEEN :F AND :T + ORDER BY SCORE) + WHERE ROWNUM = 1", new { KEY = key, F = fromScore, T = toScore })); } @@ -422,15 +415,15 @@ public override long GetCounter(string key) throw new ArgumentNullException(nameof(key)); } - const string query = @" - SELECT SUM(S.Value) - FROM (SELECT SUM(VALUE) AS Value - FROM HF_COUNTER - WHERE KEY = :KEY - UNION ALL - SELECT VALUE as Value - FROM HF_AGGREGATED_COUNTER - WHERE KEY = :KEY) AS S"; + var query = $@" + SELECT SUM(S.Value) + FROM (SELECT SUM(VALUE) AS Value + FROM {T("Counter")} + WHERE KEY = :KEY + UNION ALL + SELECT VALUE as Value + FROM {T("AggregatedCounter")} + WHERE KEY = :KEY) AS S"; return _storage @@ -446,7 +439,7 @@ public override long GetHashCount(string key) _storage .UseConnection(connection => connection.QuerySingle( - "SELECT COUNT(ID) FROM HF_HASH WHERE KEY = :KEY", + $"SELECT COUNT(ID) FROM {T("Hash")} WHERE KEY = :KEY", new { KEY = key })); } @@ -458,7 +451,7 @@ public override TimeSpan GetHashTtl(string key) { var result = connection.QuerySingle( - "SELECT MIN(EXPIRE_AT) FROM HF_HASH WHERE KEY = :KEY", + $"SELECT MIN(EXPIRE_AT) FROM {T("Hash")} WHERE KEY = :KEY", new { KEY = key }); if (!result.HasValue) @@ -478,7 +471,7 @@ public override long GetListCount(string key) _storage .UseConnection(connection => connection.QuerySingle( - "SELECT COUNT(ID) FROM HF_LIST WHERE KEY = :KEY", + $"SELECT COUNT(ID) FROM {T("List")} WHERE KEY = :KEY", new { KEY = key })); } @@ -489,7 +482,7 @@ public override TimeSpan GetListTtl(string key) return _storage.UseConnection(connection => { var result = connection.QuerySingle( - "SELECT MIN(EXPIRE_AT) FROM HF_LIST WHERE KEY = :KEY", + $"SELECT MIN(EXPIRE_AT) FROM {T("List")} WHERE KEY = :KEY", new { KEY = key }); if (!result.HasValue) @@ -510,7 +503,7 @@ public override string GetValueFromHash(string key, string name) _storage .UseConnection(connection => connection.QuerySingleOrDefault( - "SELECT VALUE AS Value FROM HF_HASH WHERE KEY = :KEY and FIELD = :FIELD", + $"SELECT VALUE AS Value FROM {T("Hash")} WHERE KEY = :KEY and FIELD = :FIELD", new { KEY = key, FIELD = name })); } @@ -521,13 +514,13 @@ public override List GetRangeFromList(string key, int startingFrom, int throw new ArgumentNullException(nameof(key)); } - const string query = @" -SELECT VALUE as Value - FROM (SELECT VALUE, RANK () OVER (ORDER BY ID DESC) AS RANK - FROM HF_LIST - WHERE KEY = :KEY) - WHERE RANK BETWEEN :S AND :E -"; + var query = $@" + SELECT VALUE as Value + FROM (SELECT VALUE, RANK () OVER (ORDER BY ID DESC) AS RANK + FROM {T("List")} + WHERE KEY = :KEY) + WHERE RANK BETWEEN :S AND :E"; + return _storage .UseConnection(connection => @@ -540,11 +533,11 @@ public override List GetAllItemsFromList(string key) { if (key == null) throw new ArgumentNullException(nameof(key)); - const string query = @" - SELECT VALUE AS Value - FROM HF_LIST - WHERE KEY = :KEY - ORDER BY ID DESC"; + var query = $@" + SELECT VALUE AS Value + FROM {T("List")} + WHERE KEY = :KEY + ORDER BY ID DESC"; return _storage.UseConnection(connection => connection.Query(query, new { KEY = key }).ToList()); } @@ -555,7 +548,9 @@ public override TimeSpan GetSetTtl(string key) return _storage.UseConnection(connection => { - var result = connection.QuerySingle("SELECT MIN(EXPIRE_AT) FROM HF_SET WHERE KEY = :KEY", new { KEY = key }); + var result = connection.QuerySingle( + $"SELECT MIN(EXPIRE_AT) FROM {T("Set")} WHERE KEY = :KEY", + new { KEY = key }); if (!result.HasValue) { @@ -587,16 +582,14 @@ public override void SetRangeInHash(string key, IEnumerable GetAllEntriesFromHash(string key) return _storage.UseConnection(connection => { var result = connection.Query( - "SELECT FIELD AS Field, VALUE AS Value FROM HF_HASH WHERE KEY = :KEY", + $"SELECT FIELD AS Field, VALUE AS Value FROM {T("Hash")} WHERE KEY = :KEY", new { KEY = key }) .ToDictionary(x => x.Field, x => x.Value); diff --git a/Hangfire.Oracle/OracleStorageOptions.cs b/Hangfire.Oracle/OracleStorageOptions.cs index 79c64a1..ddc7068 100644 --- a/Hangfire.Oracle/OracleStorageOptions.cs +++ b/Hangfire.Oracle/OracleStorageOptions.cs @@ -1,3 +1,4 @@ +using Hangfire.Oracle.Core.Configuration; using System; using System.Data; @@ -17,6 +18,7 @@ public OracleStorageOptions() DashboardJobListLimit = 50000; TransactionTimeout = TimeSpan.FromMinutes(1); InvisibilityTimeout = TimeSpan.FromMinutes(30); + InstanceConfiguration = null; } public IsolationLevel? TransactionIsolationLevel { get; set; } @@ -42,15 +44,14 @@ public TimeSpan QueuePollInterval } public bool PrepareSchemaIfNecessary { get; set; } - public TimeSpan JobExpirationCheckInterval { get; set; } public TimeSpan CountersAggregateInterval { get; set; } - public int? DashboardJobListLimit { get; set; } public TimeSpan TransactionTimeout { get; set; } + [Obsolete("Does not make sense anymore. Background jobs re-queued instantly even after ungraceful shutdown now. Will be removed in 2.0.0.")] public TimeSpan InvisibilityTimeout { get; set; } - public string SchemaName { get; set; } + public HangfireConfiguration InstanceConfiguration { get; set; } } } \ No newline at end of file diff --git a/Hangfire.Oracle/OracleWriteOnlyTransaction.cs b/Hangfire.Oracle/OracleWriteOnlyTransaction.cs index e4e9d48..f0cbe57 100644 --- a/Hangfire.Oracle/OracleWriteOnlyTransaction.cs +++ b/Hangfire.Oracle/OracleWriteOnlyTransaction.cs @@ -2,10 +2,8 @@ using System.Collections.Generic; using System.Data; using System.Linq; - using Dapper; using Dapper.Oracle; - using Hangfire.Common; using Hangfire.Logging; using Hangfire.Oracle.Core.Entities; @@ -19,7 +17,6 @@ internal class OracleWriteOnlyTransaction : JobStorageTransaction private static readonly ILog Logger = LogProvider.GetLogger(typeof(OracleWriteOnlyTransaction)); private readonly OracleStorage _storage; - private readonly Queue> _commandQueue = new Queue>(); public OracleWriteOnlyTransaction(OracleStorage storage) @@ -27,6 +24,9 @@ public OracleWriteOnlyTransaction(OracleStorage storage) _storage = storage ?? throw new ArgumentNullException(nameof(storage)); } + private string T(string logicalName) => _storage.TableNameProvider.GetTableName(logicalName); + private string GetPrimarySequence() => _storage.TableNameProvider.GetPrimarySequenceName(); + public override void ExpireJob(string jobId, TimeSpan expireIn) { Logger.TraceFormat("ExpireJob jobId={0}", jobId); @@ -35,7 +35,7 @@ public override void ExpireJob(string jobId, TimeSpan expireIn) QueueCommand(x => x.Execute( - "UPDATE HF_JOB SET EXPIRE_AT = :EXPIRE_AT WHERE ID = :ID", + $"UPDATE {T("Job")} SET EXPIRE_AT = :EXPIRE_AT WHERE ID = :ID", new { EXPIRE_AT = DateTime.UtcNow.Add(expireIn), ID = jobId })); } @@ -45,7 +45,7 @@ public override void PersistJob(string jobId) AcquireJobLock(); - QueueCommand(x => x.Execute("UPDATE HF_JOB SET EXPIRE_AT = NULL WHERE ID = :ID", new { ID = jobId })); + QueueCommand(x => x.Execute($"UPDATE {T("Job")} SET EXPIRE_AT = NULL WHERE ID = :ID", new { ID = jobId })); } public override void SetJobState(string jobId, IState state) @@ -55,7 +55,7 @@ public override void SetJobState(string jobId, IState state) AcquireStateLock(); AcquireJobLock(); - var stateId = _storage.UseConnection(connection => connection.GetNextId()); + var stateId = _storage.UseConnection(connection => connection.GetNextId(GetPrimarySequence())); var oracleDynamicParameters = new OracleDynamicParameters(); oracleDynamicParameters.AddDynamicParams(new @@ -70,14 +70,12 @@ public override void SetJobState(string jobId, IState state) oracleDynamicParameters.Add("DATA", SerializationHelper.Serialize(state.SerializeData(), SerializationOption.User), OracleMappingType.NClob, ParameterDirection.Input); QueueCommand(x => x.Execute( - @" -BEGIN - INSERT INTO HF_JOB_STATE (ID, JOB_ID, NAME, REASON, CREATED_AT, DATA) - VALUES (:STATE_ID, :JOB_ID, :NAME, :REASON, :CREATED_AT, :DATA); - - UPDATE HF_JOB SET STATE_ID = :STATE_ID, STATE_NAME = :NAME WHERE ID = :ID; -END; -", + $@"BEGIN + INSERT INTO {T("JobState")} (ID, JOB_ID, NAME, REASON, CREATED_AT, DATA) + VALUES (:STATE_ID, :JOB_ID, :NAME, :REASON, :CREATED_AT, :DATA); + + UPDATE {T("Job")} SET STATE_ID = :STATE_ID, STATE_NAME = :NAME WHERE ID = :ID; + END;", oracleDynamicParameters)); } @@ -98,8 +96,9 @@ public override void AddJobState(string jobId, IState state) oracleDynamicParameters.Add("DATA", SerializationHelper.Serialize(state.SerializeData(), SerializationOption.User), OracleMappingType.NClob, ParameterDirection.Input); QueueCommand(x => x.Execute( - " INSERT INTO HF_JOB_STATE (ID, JOB_ID, NAME, REASON, CREATED_AT, DATA) " + - " VALUES (HF_SEQUENCE.NEXTVAL, :JOB_ID, :NAME, :REASON, :CREATED_AT, :DATA)", oracleDynamicParameters)); + $@"INSERT INTO {T("JobState")} (ID, JOB_ID, NAME, REASON, CREATED_AT, DATA) + VALUES ({GetPrimarySequence()}.NEXTVAL, :JOB_ID, :NAME, :REASON, :CREATED_AT, :DATA)", + oracleDynamicParameters)); } public override void AddToQueue(string queue, string jobId) @@ -118,11 +117,11 @@ public override void IncrementCounter(string key) AcquireCounterLock(); - QueueCommand(x => x.Execute("INSERT INTO HF_COUNTER (ID, KEY, VALUE) values (HF_SEQUENCE.NEXTVAL, :KEY, :VALUE)", - new { KEY = key, VALUE = +1 })); + QueueCommand(x => x.Execute( + $"INSERT INTO {T("Counter")} (ID, KEY, VALUE) values ({GetPrimarySequence()}.NEXTVAL, :KEY, :VALUE)", + new { KEY = key, VALUE = +1 })); } - public override void IncrementCounter(string key, TimeSpan expireIn) { Logger.TraceFormat("IncrementCounter key={0}, expireIn={1}", key, expireIn); @@ -131,7 +130,7 @@ public override void IncrementCounter(string key, TimeSpan expireIn) QueueCommand(x => x.Execute( - "INSERT INTO HF_COUNTER (ID, KEY, VALUE, EXPIRE_AT) values (HF_SEQUENCE.NEXTVAL, :KEY, :VALUE, :EXPIRE_AT)", + $"INSERT INTO {T("Counter")} (ID, KEY, VALUE, EXPIRE_AT) values ({GetPrimarySequence()}.NEXTVAL, :KEY, :VALUE, :EXPIRE_AT)", new { KEY = key, VALUE = +1, EXPIRE_AT = DateTime.UtcNow.Add(expireIn) })); } @@ -143,7 +142,7 @@ public override void DecrementCounter(string key) QueueCommand(x => x.Execute( - "INSERT INTO HF_COUNTER (ID, KEY, VALUE) values (HF_SEQUENCE.NEXTVAL, :KEY, :VALUE)", + $"INSERT INTO {T("Counter")} (ID, KEY, VALUE) values ({GetPrimarySequence()}.NEXTVAL, :KEY, :VALUE)", new { KEY = key, VALUE = -1 })); } @@ -154,7 +153,7 @@ public override void DecrementCounter(string key, TimeSpan expireIn) AcquireCounterLock(); QueueCommand(x => x.Execute( - "INSERT INTO HF_COUNTER (ID, KEY, VALUE, EXPIRE_AT) values (HF_SEQUENCE.NEXTVAL, :KEY, :VALUE, :EXPIRE_AT)", + $"INSERT INTO {T("Counter")} (ID, KEY, VALUE, EXPIRE_AT) values ({GetPrimarySequence()}.NEXTVAL, :KEY, :VALUE, :EXPIRE_AT)", new { KEY = key, VALUE = -1, EXPIRE_AT = DateTime.UtcNow.Add(expireIn) })); } @@ -170,16 +169,14 @@ public override void AddToSet(string key, string value, double score) AcquireSetLock(); QueueCommand(x => x.Execute( - @" - MERGE INTO HF_SET H - USING (SELECT 1 FROM DUAL) SRC - ON (H.KEY = :KEY AND H.VALUE = :VALUE) - WHEN MATCHED THEN - UPDATE SET SCORE = :SCORE - WHEN NOT MATCHED THEN - INSERT (ID, KEY, VALUE, SCORE) - VALUES (HF_SEQUENCE.NEXTVAL, :KEY, :VALUE, :SCORE) -", + $@"MERGE INTO {T("Set")} H + USING (SELECT 1 FROM DUAL) SRC + ON (H.KEY = :KEY AND H.VALUE = :VALUE) + WHEN MATCHED THEN + UPDATE SET SCORE = :SCORE + WHEN NOT MATCHED THEN + INSERT (ID, KEY, VALUE, SCORE) + VALUES ({GetPrimarySequence()}.NEXTVAL, :KEY, :VALUE, :SCORE)", new { KEY = key, VALUE = value, SCORE = score })); } @@ -193,17 +190,18 @@ public override void AddRangeToSet(string key, IList items) AcquireSetLock(); QueueCommand(x => x.Execute( - "INSERT INTO HF_SET (ID, KEY, VALUE, SCORE) VALUES (HF_SEQUENCE.NEXTVAL, :KEY, :VALUE, 0.0)", + $"INSERT INTO {T("Set")} (ID, KEY, VALUE, SCORE) VALUES ({GetPrimarySequence()}.NEXTVAL, :KEY, :VALUE, 0.0)", items.Select(value => new { KEY = key, VALUE = value }).ToList())); } - public override void RemoveFromSet(string key, string value) { Logger.TraceFormat("RemoveFromSet key={0} value={1}", key, value); AcquireSetLock(); - QueueCommand(x => x.Execute("DELETE FROM HF_SET WHERE KEY = :KEY AND VALUE = :VALUE", new { KEY = key, VALUE = value })); + QueueCommand(x => x.Execute( + $"DELETE FROM {T("Set")} WHERE KEY = :KEY AND VALUE = :VALUE", + new { KEY = key, VALUE = value })); } public override void ExpireSet(string key, TimeSpan expireIn) @@ -215,7 +213,7 @@ public override void ExpireSet(string key, TimeSpan expireIn) AcquireSetLock(); QueueCommand(x => x.Execute( - "UPDATE HF_SET SET EXPIRE_AT = :EXPIRE_AT WHERE KEY = :KEY", + $"UPDATE {T("Set")} SET EXPIRE_AT = :EXPIRE_AT WHERE KEY = :KEY", new { KEY = key, EXPIRE_AT = DateTime.UtcNow.Add(expireIn) })); } @@ -229,7 +227,9 @@ public override void InsertToList(string key, string value) oracleDynamicParameters.Add("KEY", key); oracleDynamicParameters.Add("VALUE", value, OracleMappingType.NClob, ParameterDirection.Input); - QueueCommand(x => x.Execute("INSERT INTO HF_LIST (ID, KEY, VALUE) VALUES (HF_SEQUENCE.NEXTVAL, :KEY, :VALUE)", oracleDynamicParameters)); + QueueCommand(x => x.Execute( + $"INSERT INTO {T("List")} (ID, KEY, VALUE) VALUES ({GetPrimarySequence()}.NEXTVAL, :KEY, :VALUE)", + oracleDynamicParameters)); } public override void ExpireList(string key, TimeSpan expireIn) @@ -244,7 +244,7 @@ public override void ExpireList(string key, TimeSpan expireIn) AcquireListLock(); QueueCommand(x => x.Execute( - "UPDATE HF_LIST SET EXPIRE_AT = :EXPIRE_AT WHERE KEY = :KEY", + $"UPDATE {T("List")} SET EXPIRE_AT = :EXPIRE_AT WHERE KEY = :KEY", new { KEY = key, EXPIRE_AT = DateTime.UtcNow.Add(expireIn) })); } @@ -254,7 +254,7 @@ public override void RemoveFromList(string key, string value) AcquireListLock(); QueueCommand(x => x.Execute( - "DELETE FROM HF_LIST WHERE KEY = :KEY AND VALUE = :VALUE", + $"DELETE FROM {T("List")} WHERE KEY = :KEY AND VALUE = :VALUE", new { KEY = key, VALUE = value })); } @@ -264,14 +264,16 @@ public override void TrimList(string key, int keepStartingFrom, int keepEndingAt AcquireListLock(); QueueCommand(x => x.Execute( - @" -delete lst -from List lst - inner join (SELECT tmp.Id, @rownum := @rownum + 1 AS rankvalue - FROM List tmp, - (SELECT @rownum := 0) r ) ranked on ranked.Id = lst.Id -where lst.Key = @key - and ranked.rankvalue not between @start and @end", + $@"DELETE FROM {T("List")} lst + WHERE lst.ID IN ( + SELECT tmp.ID + FROM ( + SELECT ID, ROW_NUMBER() OVER (ORDER BY ID) AS rankvalue + FROM {T("List")} + WHERE KEY = :key + ) tmp + WHERE tmp.rankvalue NOT BETWEEN :start AND :end + )", new { key, start = keepStartingFrom + 1, end = keepEndingAt + 1 })); } @@ -284,7 +286,8 @@ public override void PersistHash(string key) AcquireHashLock(); QueueCommand(x => x.Execute( - "UPDATE HF_HASH SET EXPIRE_AT = NULL WHERE KEY = :KEY", new { KEY = key })); + $"UPDATE {T("Hash")} SET EXPIRE_AT = NULL WHERE KEY = :KEY", + new { KEY = key })); } public override void PersistSet(string key) @@ -294,7 +297,9 @@ public override void PersistSet(string key) if (key == null) throw new ArgumentNullException(nameof(key)); AcquireSetLock(); - QueueCommand(x => x.Execute("UPDATE HF_SET SET EXPIRE_AT = NULL WHERE KEY = :KEY", new { KEY = key })); + QueueCommand(x => x.Execute( + $"UPDATE {T("Set")} SET EXPIRE_AT = NULL WHERE KEY = :KEY", + new { KEY = key })); } public override void RemoveSet(string key) @@ -304,7 +309,9 @@ public override void RemoveSet(string key) if (key == null) throw new ArgumentNullException(nameof(key)); AcquireSetLock(); - QueueCommand(x => x.Execute("DELETE FROM HF_SET WHERE KEY = :KEY", new { KEY = key })); + QueueCommand(x => x.Execute( + $"DELETE FROM {T("Set")} WHERE KEY = :KEY", + new { KEY = key })); } public override void PersistList(string key) @@ -316,7 +323,8 @@ public override void PersistList(string key) AcquireListLock(); QueueCommand(x => x.Execute( - "UPDATE HF_LIST SET EXPIRE_AT = NULL WHERE KEY = :KEY", new { KEY = key })); + $"UPDATE {T("List")} SET EXPIRE_AT = NULL WHERE KEY = :KEY", + new { KEY = key })); } public override void SetRangeInHash(string key, IEnumerable> keyValuePairs) @@ -336,16 +344,14 @@ public override void SetRangeInHash(string key, IEnumerable x.Execute( - @" - MERGE INTO HF_HASH H - USING (SELECT 1 FROM DUAL) SRC - ON (H.KEY = :KEY AND H.FIELD = :FIELD) - WHEN MATCHED THEN - UPDATE SET VALUE = :VALUE - WHEN NOT MATCHED THEN - INSERT (ID, KEY, VALUE, FIELD) - VALUES (HF_SEQUENCE.NEXTVAL, :KEY, :VALUE, :FIELD) -", + $@"MERGE INTO {T("Hash")} H + USING (SELECT 1 FROM DUAL) SRC + ON (H.KEY = :KEY AND H.FIELD = :FIELD) + WHEN MATCHED THEN + UPDATE SET VALUE = :VALUE + WHEN NOT MATCHED THEN + INSERT (ID, KEY, VALUE, FIELD) + VALUES ({GetPrimarySequence()}.NEXTVAL, :KEY, :VALUE, :FIELD)", keyValuePairs.Select(y => new { KEY = key, FIELD = y.Key, VALUE = y.Value }))); } @@ -361,7 +367,7 @@ public override void ExpireHash(string key, TimeSpan expireIn) AcquireHashLock(); QueueCommand(x => x.Execute( - "UPDATE HF_HASH SET EXPIRE_AT = :EXPIRE_AT WHERE KEY = :KEY", + $"UPDATE {T("Hash")} SET EXPIRE_AT = :EXPIRE_AT WHERE KEY = :KEY", new { KEY = key, EXPIRE_AT = DateTime.UtcNow.Add(expireIn) })); } @@ -375,7 +381,9 @@ public override void RemoveHash(string key) } AcquireHashLock(); - QueueCommand(x => x.Execute("DELETE FROM HF_HASH WHERE KEY = :KEY", new { KEY = key })); + QueueCommand(x => x.Execute( + $"DELETE FROM {T("Hash")} WHERE KEY = :KEY", + new { KEY = key })); } public override void Commit() @@ -423,6 +431,7 @@ private void AcquireCounterLock() { AcquireLock("Counter"); } + private void AcquireLock(string resource) { } diff --git a/Hangfire.Oracle/hangfire-config.json b/Hangfire.Oracle/hangfire-config.json new file mode 100644 index 0000000..7628b0f --- /dev/null +++ b/Hangfire.Oracle/hangfire-config.json @@ -0,0 +1,21 @@ +{ + "SchemaName": "AHA", + "InstanceName": "AHA", + "Sequence": { + "PrimarySequenceName": "AHA_SEQUENCE", + "JobIdSequenceName": "AHA_JOB_ID_SEQ" + }, + "Tables": { + "Job": "AHA_JOB", + "JobParameter": "AHA_JOB_PARAMETER", + "JobQueue": "AHA_JOB_QUEUE", + "JobState": "AHA_JOB_STATE", + "Server": "AHA_SERVER", + "Set": "AHA_SET", + "List": "AHA_LIST", + "Hash": "AHA_HASH", + "Counter": "AHA_COUNTER", + "AggregatedCounter": "AHA_AGGREGATED_COUNTER", + "DistributedLock": "AHA_DISTRIBUTED_LOCK" + } +} \ No newline at end of file