diff --git a/CmdScale.EntityFrameworkCore.TimescaleDB.Design/TimescaleCSharpModelGenerator.cs b/CmdScale.EntityFrameworkCore.TimescaleDB.Design/TimescaleCSharpModelGenerator.cs index 408897c..1edd468 100644 --- a/CmdScale.EntityFrameworkCore.TimescaleDB.Design/TimescaleCSharpModelGenerator.cs +++ b/CmdScale.EntityFrameworkCore.TimescaleDB.Design/TimescaleCSharpModelGenerator.cs @@ -104,7 +104,7 @@ public override DatabaseModel Create(DbConnection connection, DatabaseModelFacto private static Dictionary<(string, string), HypertableInfo> GetHypertables(DbConnection connection) { - bool wasOpen = connection.State == System.Data.ConnectionState.Open; + bool wasOpen = connection.State == ConnectionState.Open; if (!wasOpen) { connection.Open(); diff --git a/CmdScale.EntityFrameworkCore.TimescaleDB.Example.DataAccess/TimescaleContext.cs b/CmdScale.EntityFrameworkCore.TimescaleDB.Example.DataAccess/TimescaleContext.cs index 89eb315..e7bd612 100644 --- a/CmdScale.EntityFrameworkCore.TimescaleDB.Example.DataAccess/TimescaleContext.cs +++ b/CmdScale.EntityFrameworkCore.TimescaleDB.Example.DataAccess/TimescaleContext.cs @@ -16,6 +16,7 @@ protected override void OnModelCreating(ModelBuilder modelBuilder) { base.OnModelCreating(modelBuilder); modelBuilder.ApplyConfigurationsFromAssembly(typeof(TimescaleContext).Assembly); + modelBuilder.HasDefaultSchema("custom_schema"); } } } \ No newline at end of file diff --git a/CmdScale.EntityFrameworkCore.TimescaleDB.Example/Program.cs b/CmdScale.EntityFrameworkCore.TimescaleDB.Example/Program.cs index 84a3735..efb6d19 100644 --- a/CmdScale.EntityFrameworkCore.TimescaleDB.Example/Program.cs +++ b/CmdScale.EntityFrameworkCore.TimescaleDB.Example/Program.cs @@ -13,6 +13,26 @@ IHost host = builder.Build(); +#if false +// --- Code to run Database.EnsureCreatedAsync() --- +// NOTE: Set the #if to false to disable this block or to true to enable it. +using (IServiceScope scope = host.Services.CreateScope()) +{ + IServiceProvider services = scope.ServiceProvider; + try + { + TimescaleContext context = services.GetRequiredService(); + Console.WriteLine("Applying Database.EnsureCreatedAsync()..."); + await context.Database.EnsureCreatedAsync(); + Console.WriteLine("Database setup complete."); + } + catch (Exception ex) + { + Console.WriteLine($"An error occurred while creating the database: {ex.Message}"); + } +} +#endif + Console.WriteLine("TimescaleDB EF Core Demo"); Console.WriteLine("------------------------------------"); Console.WriteLine("Run 'dotnet ef migrations add ' to generate a migration."); diff --git a/CmdScale.EntityFrameworkCore.TimescaleDB.Tests/Generators/HypertableOperationGeneratorTests.cs b/CmdScale.EntityFrameworkCore.TimescaleDB.Tests/Generators/HypertableOperationGeneratorTests.cs index 0a5fdca..fe141d7 100644 --- a/CmdScale.EntityFrameworkCore.TimescaleDB.Tests/Generators/HypertableOperationGeneratorTests.cs +++ b/CmdScale.EntityFrameworkCore.TimescaleDB.Tests/Generators/HypertableOperationGeneratorTests.cs @@ -30,11 +30,12 @@ public void Generate_Create_with_minimal_details_generates_correct_sql() CreateHypertableOperation operation = new() { TableName = "MinimalTable", + Schema = "public", TimeColumnName = "Timestamp" }; string expected = @".Sql(@"" - SELECT create_hypertable('""""MinimalTable""""', 'Timestamp'); + SELECT create_hypertable('public.""""MinimalTable""""', 'Timestamp'); "")"; // Act @@ -51,6 +52,7 @@ public void Generate_Create_with_all_options_generates_comprehensive_sql() CreateHypertableOperation operation = new() { TableName = "FullTable", + Schema = "custom_schema", TimeColumnName = "EventTime", ChunkTimeInterval = "1 day", EnableCompression = true, @@ -62,12 +64,12 @@ public void Generate_Create_with_all_options_generates_comprehensive_sql() }; string expected = @".Sql(@"" - SELECT create_hypertable('""""FullTable""""', 'EventTime'); - SELECT set_chunk_time_interval('""""FullTable""""', INTERVAL '1 day'); - ALTER TABLE """"FullTable"""" SET (timescaledb.compress = true); + SELECT create_hypertable('custom_schema.""""FullTable""""', 'EventTime'); + SELECT set_chunk_time_interval('custom_schema.""""FullTable""""', INTERVAL '1 day'); + ALTER TABLE """"custom_schema"""".""""FullTable"""" SET (timescaledb.compress = true); SET timescaledb.enable_chunk_skipping = 'ON'; - SELECT enable_chunk_skipping('""""FullTable""""', 'DeviceId'); - SELECT add_dimension('""""FullTable""""', by_hash('LocationId', 4)); + SELECT enable_chunk_skipping('custom_schema.""""FullTable""""', 'DeviceId'); + SELECT add_dimension('custom_schema.""""FullTable""""', by_hash('LocationId', 4)); "")"; // Act @@ -84,6 +86,7 @@ public void Generate_Alter_WhenAddingChunkSkippingToUncompressedTable_ShouldAlso AlterHypertableOperation operation = new() { TableName = "Metrics", + Schema = "custom_schema", OldEnableCompression = false, OldChunkSkipColumns = [], EnableCompression = false, @@ -91,9 +94,9 @@ public void Generate_Alter_WhenAddingChunkSkippingToUncompressedTable_ShouldAlso }; string expected = @".Sql(@"" - ALTER TABLE """"Metrics"""" SET (timescaledb.compress = true); + ALTER TABLE """"custom_schema"""".""""Metrics"""" SET (timescaledb.compress = true); SET timescaledb.enable_chunk_skipping = 'ON'; - SELECT enable_chunk_skipping('""""Metrics""""', 'device_id'); + SELECT enable_chunk_skipping('custom_schema.""""Metrics""""', 'device_id'); "")"; // Act @@ -112,12 +115,13 @@ public void Generate_Alter_when_changing_compression_generates_correct_sql() AlterHypertableOperation operation = new() { TableName = "SensorData", + Schema = "public", EnableCompression = true, OldEnableCompression = false }; string expected = @".Sql(@"" - ALTER TABLE """"SensorData"""" SET (timescaledb.compress = true); + ALTER TABLE """"public"""".""""SensorData"""" SET (timescaledb.compress = true); "")"; // Act @@ -134,14 +138,15 @@ public void Generate_Alter_when_adding_and_removing_skip_columns_generates_corre AlterHypertableOperation operation = new() { TableName = "Metrics", + Schema = "metrics_schema", ChunkSkipColumns = ["host", "service"], OldChunkSkipColumns = ["host", "region"] }; string expected = @".Sql(@"" SET timescaledb.enable_chunk_skipping = 'ON'; - SELECT enable_chunk_skipping('""""Metrics""""', 'service'); - SELECT disable_chunk_skipping('""""Metrics""""', 'region'); + SELECT enable_chunk_skipping('metrics_schema.""""Metrics""""', 'service'); + SELECT disable_chunk_skipping('metrics_schema.""""Metrics""""', 'region'); "")"; // Act @@ -158,6 +163,7 @@ public void Generate_Alter_when_no_properties_change_generates_no_sql() AlterHypertableOperation operation = new() { TableName = "NoChangeTable", + Schema = "public", EnableCompression = true, OldEnableCompression = true, ChunkTimeInterval = "7 days", @@ -180,14 +186,15 @@ public void Generate_Alter_WhenRemovingLastChunkSkipColumn_ShouldDisableCompress AlterHypertableOperation operation = new() { TableName = "Logs", + Schema = "public", OldEnableCompression = false, OldChunkSkipColumns = ["trace_id"], EnableCompression = false, ChunkSkipColumns = [] }; string expected = @".Sql(@"" - ALTER TABLE """"Logs"""" SET (timescaledb.compress = false); - SELECT disable_chunk_skipping('""""Logs""""', 'trace_id'); + ALTER TABLE """"public"""".""""Logs"""" SET (timescaledb.compress = false); + SELECT disable_chunk_skipping('public.""""Logs""""', 'trace_id'); "")"; // Act diff --git a/CmdScale.EntityFrameworkCore.TimescaleDB.Tests/Generators/ReorderPolicyOperationGeneratorTests.cs b/CmdScale.EntityFrameworkCore.TimescaleDB.Tests/Generators/ReorderPolicyOperationGeneratorTests.cs index 3f74d6f..1a080f0 100644 --- a/CmdScale.EntityFrameworkCore.TimescaleDB.Tests/Generators/ReorderPolicyOperationGeneratorTests.cs +++ b/CmdScale.EntityFrameworkCore.TimescaleDB.Tests/Generators/ReorderPolicyOperationGeneratorTests.cs @@ -25,12 +25,13 @@ public void Generate_Add_with_minimal_details_creates_only_add_policy_sql() // Arrange AddReorderPolicyOperation operation = new() { + Schema = "public", TableName = "TestTable", IndexName = "IX_TestTable_Time" }; string expected = @".Sql(@"" - SELECT add_reorder_policy('""""TestTable""""', 'IX_TestTable_Time'); + SELECT add_reorder_policy('public.""""TestTable""""', 'IX_TestTable_Time'); "")"; // Act @@ -47,20 +48,21 @@ public void Generate_Add_with_non_default_schedule_creates_add_and_alter_sql() DateTime testDate = new(2025, 10, 20, 12, 30, 0, DateTimeKind.Utc); AddReorderPolicyOperation operation = new() { + Schema = "custom", TableName = "TestTable", IndexName = "IX_TestTable_Time", InitialStart = testDate, ScheduleInterval = "2 days", - MaxRuntime = "1 hour", + MaxRuntime = "1 hour", MaxRetries = 5, RetryPeriod = "10 minutes" }; string expected = @".Sql(@"" - SELECT add_reorder_policy('""""TestTable""""', 'IX_TestTable_Time', initial_start => '2025-10-20T12:30:00.0000000Z'); + SELECT add_reorder_policy('custom.""""TestTable""""', 'IX_TestTable_Time', initial_start => '2025-10-20T12:30:00.0000000Z'); SELECT alter_job(job_id, schedule_interval => INTERVAL '2 days', max_runtime => INTERVAL '1 hour', max_retries => 5, retry_period => INTERVAL '10 minutes') FROM timescaledb_information.jobs - WHERE proc_name = 'policy_reorder' AND hypertable_name = 'TestTable'; + WHERE proc_name = 'policy_reorder' AND hypertable_schema = 'custom' AND hypertable_name = 'TestTable'; "")"; // Act @@ -76,10 +78,14 @@ FROM timescaledb_information.jobs public void Generate_Drop_creates_correct_remove_policy_sql() { // Arrange - DropReorderPolicyOperation operation = new() { TableName = "TestTable" }; + DropReorderPolicyOperation operation = new() + { + Schema = "public", + TableName = "TestTable" + }; string expected = @".Sql(@"" - SELECT remove_reorder_policy('""""TestTable""""', if_exists => true); + SELECT remove_reorder_policy('public.""""TestTable""""', if_exists => true); "")"; // Act @@ -97,6 +103,7 @@ public void Generate_Alter_when_only_job_settings_change_creates_only_alter_job_ // Arrange AlterReorderPolicyOperation operation = new() { + Schema = "metrics", TableName = "TestTable", // Fundamental properties are the same IndexName = "IX_TestTable_Time", @@ -111,7 +118,7 @@ public void Generate_Alter_when_only_job_settings_change_creates_only_alter_job_ string expected = @".Sql(@"" SELECT alter_job(job_id, schedule_interval => INTERVAL '2 days') FROM timescaledb_information.jobs - WHERE proc_name = 'policy_reorder' AND hypertable_name = 'TestTable'; + WHERE proc_name = 'policy_reorder' AND hypertable_schema = 'metrics' AND hypertable_name = 'TestTable'; "")"; // Act @@ -127,6 +134,7 @@ public void Generate_Alter_when_fundamental_property_changes_creates_drop_and_ad // Arrange AlterReorderPolicyOperation operation = new() { + Schema = "logs", TableName = "TestTable", IndexName = "IX_New_Name", OldIndexName = "IX_Old_Name", @@ -135,11 +143,11 @@ public void Generate_Alter_when_fundamental_property_changes_creates_drop_and_ad }; string expected = @".Sql(@"" - SELECT remove_reorder_policy('""""TestTable""""', if_exists => true); - SELECT add_reorder_policy('""""TestTable""""', 'IX_New_Name'); + SELECT remove_reorder_policy('logs.""""TestTable""""', if_exists => true); + SELECT add_reorder_policy('logs.""""TestTable""""', 'IX_New_Name'); SELECT alter_job(job_id, schedule_interval => INTERVAL '2 days') FROM timescaledb_information.jobs - WHERE proc_name = 'policy_reorder' AND hypertable_name = 'TestTable'; + WHERE proc_name = 'policy_reorder' AND hypertable_schema = 'logs' AND hypertable_name = 'TestTable'; "")"; // Act @@ -155,6 +163,7 @@ public void Generate_Alter_when_both_fundamental_and_job_settings_change_creates // Arrange AlterReorderPolicyOperation operation = new() { + Schema = "public", TableName = "TestTable", IndexName = "IX_New_Name", OldIndexName = "IX_Old_Name", @@ -167,11 +176,11 @@ public void Generate_Alter_when_both_fundamental_and_job_settings_change_creates }; string expected = @".Sql(@"" - SELECT remove_reorder_policy('""""TestTable""""', if_exists => true); - SELECT add_reorder_policy('""""TestTable""""', 'IX_New_Name'); + SELECT remove_reorder_policy('public.""""TestTable""""', if_exists => true); + SELECT add_reorder_policy('public.""""TestTable""""', 'IX_New_Name'); SELECT alter_job(job_id, schedule_interval => INTERVAL '2 days', max_retries => 5, retry_period => INTERVAL '10 minutes') FROM timescaledb_information.jobs - WHERE proc_name = 'policy_reorder' AND hypertable_name = 'TestTable'; + WHERE proc_name = 'policy_reorder' AND hypertable_schema = 'public' AND hypertable_name = 'TestTable'; "")"; // Act @@ -181,4 +190,4 @@ FROM timescaledb_information.jobs Assert.Equal(SqlHelper.NormalizeSql(expected), SqlHelper.NormalizeSql(result)); } } -} +} \ No newline at end of file diff --git a/CmdScale.EntityFrameworkCore.TimescaleDB.Tests/Generators/SqlBuilderHelperTests.cs b/CmdScale.EntityFrameworkCore.TimescaleDB.Tests/Generators/SqlBuilderHelperTests.cs index 7e9cf42..aaec120 100644 --- a/CmdScale.EntityFrameworkCore.TimescaleDB.Tests/Generators/SqlBuilderHelperTests.cs +++ b/CmdScale.EntityFrameworkCore.TimescaleDB.Tests/Generators/SqlBuilderHelperTests.cs @@ -19,7 +19,7 @@ public void Regclass_Runtime_ReturnsCorrectlyQuotedString() // Arrange SqlBuilderHelper helper = new(quoteString: "\""); string tableName = "MyTable"; - string expected = "'\"MyTable\"'"; + string expected = "'public.\"MyTable\"'"; // Act string result = helper.Regclass(tableName); @@ -34,7 +34,7 @@ public void QualifiedIdentifier_Runtime_ReturnsCorrectlyQuotedString() // Arrange SqlBuilderHelper helper = new(quoteString: "\""); string tableName = "MyTable"; - string expected = "\"MyTable\""; + string expected = "\"public\".\"MyTable\""; // Act string result = helper.QualifiedIdentifier(tableName); @@ -49,7 +49,7 @@ public void Regclass_DesignTime_ReturnsCorrectlyEscapedQuotedString() // Arrange SqlBuilderHelper helper = new(quoteString: "\"\""); string tableName = "MyTable"; - string expected = "'\"\"MyTable\"\"'"; + string expected = "'public.\"\"MyTable\"\"'"; // Act string result = helper.Regclass(tableName); @@ -64,7 +64,7 @@ public void QualifiedIdentifier_DesignTime_ReturnsCorrectlyEscapedQuotedString() // Arrange SqlBuilderHelper helper = new(quoteString: "\"\""); string tableName = "MyTable"; - string expected = "\"\"MyTable\"\""; + string expected = "\"\"public\"\".\"\"MyTable\"\""; // Act string result = helper.QualifiedIdentifier(tableName); diff --git a/CmdScale.EntityFrameworkCore.TimescaleDB/DefaultValues.cs b/CmdScale.EntityFrameworkCore.TimescaleDB/DefaultValues.cs index 8b6d3e4..74a6255 100644 --- a/CmdScale.EntityFrameworkCore.TimescaleDB/DefaultValues.cs +++ b/CmdScale.EntityFrameworkCore.TimescaleDB/DefaultValues.cs @@ -5,6 +5,7 @@ /// public static class DefaultValues { + public const string DefaultSchema = "public"; public const string ChunkTimeInterval = "7 days"; public const long ChunkTimeIntervalLong = 604_800_000_000L; public const string ReorderPolicyScheduleInterval = "1 day"; diff --git a/CmdScale.EntityFrameworkCore.TimescaleDB/Generators/HypertableOperationGenerator.cs b/CmdScale.EntityFrameworkCore.TimescaleDB/Generators/HypertableOperationGenerator.cs index 8b81bce..1582f73 100644 --- a/CmdScale.EntityFrameworkCore.TimescaleDB/Generators/HypertableOperationGenerator.cs +++ b/CmdScale.EntityFrameworkCore.TimescaleDB/Generators/HypertableOperationGenerator.cs @@ -21,9 +21,12 @@ public HypertableOperationGenerator(bool isDesignTime = false) public List Generate(CreateHypertableOperation operation) { + string qualifiedTableName = sqlHelper.Regclass(operation.TableName, operation.Schema); + string qualifiedIdentifier = sqlHelper.QualifiedIdentifier(operation.TableName, operation.Schema); + List statements = [ - $"SELECT create_hypertable({sqlHelper.Regclass(operation.TableName)}, '{operation.TimeColumnName}');" + $"SELECT create_hypertable({qualifiedTableName}, '{operation.TimeColumnName}');" ]; // ChunkTimeInterval @@ -33,12 +36,12 @@ public List Generate(CreateHypertableOperation operation) if (long.TryParse(operation.ChunkTimeInterval, out _)) { // If it's a number, don't wrap it in quotes. - statements.Add($"SELECT set_chunk_time_interval({sqlHelper.Regclass(operation.TableName)}, {operation.ChunkTimeInterval}::bigint);"); + statements.Add($"SELECT set_chunk_time_interval({qualifiedTableName}, {operation.ChunkTimeInterval}::bigint);"); } else { // If it's a string like '7 days', wrap it in quotes. - statements.Add($"SELECT set_chunk_time_interval({sqlHelper.Regclass(operation.TableName)}, INTERVAL '{operation.ChunkTimeInterval}');"); + statements.Add($"SELECT set_chunk_time_interval({qualifiedTableName}, INTERVAL '{operation.ChunkTimeInterval}');"); } } @@ -46,7 +49,7 @@ public List Generate(CreateHypertableOperation operation) if (operation.EnableCompression || operation.ChunkSkipColumns?.Count > 0) { bool enableCompression = operation.EnableCompression || operation.ChunkSkipColumns != null && operation.ChunkSkipColumns.Count > 0; - statements.Add($"ALTER TABLE {sqlHelper.QualifiedIdentifier(operation.TableName)} SET (timescaledb.compress = {enableCompression.ToString().ToLower()});"); + statements.Add($"ALTER TABLE {qualifiedIdentifier} SET (timescaledb.compress = {enableCompression.ToString().ToLower()});"); } // ChunkSkipColumns @@ -56,7 +59,7 @@ public List Generate(CreateHypertableOperation operation) foreach (string column in operation.ChunkSkipColumns) { - statements.Add($"SELECT enable_chunk_skipping({sqlHelper.Regclass(operation.TableName)}, '{column}');"); + statements.Add($"SELECT enable_chunk_skipping({qualifiedTableName}, '{column}');"); } } @@ -67,11 +70,11 @@ public List Generate(CreateHypertableOperation operation) { if (dimension.Type == EDimensionType.Range) { - statements.Add($"SELECT add_dimension({sqlHelper.Regclass(operation.TableName)}, by_range('{dimension.ColumnName}', INTERVAL '{dimension.Interval}'));"); + statements.Add($"SELECT add_dimension({qualifiedTableName}, by_range('{dimension.ColumnName}', INTERVAL '{dimension.Interval}'));"); } else if (dimension.Type == EDimensionType.Hash) { - statements.Add($"SELECT add_dimension({sqlHelper.Regclass(operation.TableName)}, by_hash('{dimension.ColumnName}', {dimension.NumberOfPartitions}));"); + statements.Add($"SELECT add_dimension({qualifiedTableName}, by_hash('{dimension.ColumnName}', {dimension.NumberOfPartitions}));"); } } } @@ -81,6 +84,9 @@ public List Generate(CreateHypertableOperation operation) public List Generate(AlterHypertableOperation operation) { + string qualifiedTableName = sqlHelper.Regclass(operation.TableName, operation.Schema); + string qualifiedIdentifier = sqlHelper.QualifiedIdentifier(operation.TableName, operation.Schema); + List statements = []; // Check for ChunkTimeInterval change @@ -90,12 +96,12 @@ public List Generate(AlterHypertableOperation operation) if (long.TryParse(operation.ChunkTimeInterval, out _)) { // If it's a number, don't wrap it in quotes. - statements.Add($"SELECT set_chunk_time_interval({sqlHelper.Regclass(operation.TableName)}, {operation.ChunkTimeInterval}::bigint);"); + statements.Add($"SELECT set_chunk_time_interval({qualifiedTableName}, {operation.ChunkTimeInterval}::bigint);"); } else { // If it's a string like '7 days', wrap it in quotes. - statements.Add($"SELECT set_chunk_time_interval({sqlHelper.Regclass(operation.TableName)}, INTERVAL '{operation.ChunkTimeInterval}');"); + statements.Add($"SELECT set_chunk_time_interval({qualifiedTableName}, INTERVAL '{operation.ChunkTimeInterval}');"); } } @@ -106,7 +112,7 @@ public List Generate(AlterHypertableOperation operation) if (newCompressionState != oldCompressionState) { string compressionValue = newCompressionState.ToString().ToLower(); - statements.Add($"ALTER TABLE {sqlHelper.QualifiedIdentifier(operation.TableName)} SET (timescaledb.compress = {compressionValue});"); + statements.Add($"ALTER TABLE {qualifiedIdentifier} SET (timescaledb.compress = {compressionValue});"); } // Handle ChunkSkipColumns @@ -120,7 +126,7 @@ public List Generate(AlterHypertableOperation operation) foreach (string column in addedColumns) { - statements.Add($"SELECT enable_chunk_skipping({sqlHelper.Regclass(operation.TableName)}, '{column}');"); + statements.Add($"SELECT enable_chunk_skipping({qualifiedTableName}, '{column}');"); } } @@ -129,7 +135,7 @@ public List Generate(AlterHypertableOperation operation) { foreach (string column in removedColumns) { - statements.Add($"SELECT disable_chunk_skipping({sqlHelper.Regclass(operation.TableName)}, '{column}');"); + statements.Add($"SELECT disable_chunk_skipping({qualifiedTableName}, '{column}');"); } } diff --git a/CmdScale.EntityFrameworkCore.TimescaleDB/Generators/ReorderPolicyOperationGenerator.cs b/CmdScale.EntityFrameworkCore.TimescaleDB/Generators/ReorderPolicyOperationGenerator.cs index b69fd5a..294ec4f 100644 --- a/CmdScale.EntityFrameworkCore.TimescaleDB/Generators/ReorderPolicyOperationGenerator.cs +++ b/CmdScale.EntityFrameworkCore.TimescaleDB/Generators/ReorderPolicyOperationGenerator.cs @@ -22,13 +22,13 @@ public List Generate(AddReorderPolicyOperation operation) { List statements = [ - BuildAddReorderPolicySql(operation.TableName, operation.IndexName, operation.InitialStart) + BuildAddReorderPolicySql(operation.TableName, operation.Schema, operation.IndexName, operation.InitialStart) ]; List alterJobClauses = BuildAlterJobClauses(operation); if (alterJobClauses.Count != 0) { - statements.Add(BuildAlterJobSql(operation.TableName, alterJobClauses)); + statements.Add(BuildAlterJobSql(operation.TableName, operation.Schema, alterJobClauses)); } return statements; @@ -36,13 +36,15 @@ public List Generate(AddReorderPolicyOperation operation) public List Generate(AlterReorderPolicyOperation operation) { + string qualifiedTableName = sqlHelper.Regclass(operation.TableName, operation.Schema); + List statements = []; bool needsRecreation = operation.IndexName != operation.OldIndexName || operation.InitialStart != operation.OldInitialStart; if (needsRecreation) { - statements.Add($"SELECT remove_reorder_policy({sqlHelper.Regclass(operation.TableName)}, if_exists => true);"); - statements.Add(BuildAddReorderPolicySql(operation.TableName, operation.IndexName, operation.InitialStart)); + statements.Add($"SELECT remove_reorder_policy({qualifiedTableName}, if_exists => true);"); + statements.Add(BuildAddReorderPolicySql(operation.TableName, operation.Schema, operation.IndexName, operation.InitialStart)); // Create a temporary "add" operation representing the final desired state to ensure existing settings are reapplied. AddReorderPolicyOperation finalStateOperation = new() @@ -59,7 +61,7 @@ public List Generate(AlterReorderPolicyOperation operation) List finalStateClauses = BuildAlterJobClauses(finalStateOperation); if (finalStateClauses.Count != 0) { - statements.Add(BuildAlterJobSql(operation.TableName, finalStateClauses)); + statements.Add(BuildAlterJobSql(operation.TableName, operation.Schema, finalStateClauses)); } } else @@ -67,7 +69,7 @@ public List Generate(AlterReorderPolicyOperation operation) List changedClauses = BuildAlterJobClauses(operation); if (changedClauses.Count != 0) { - statements.Add(BuildAlterJobSql(operation.TableName, changedClauses)); + statements.Add(BuildAlterJobSql(operation.TableName, operation.Schema, changedClauses)); } } @@ -76,9 +78,11 @@ public List Generate(AlterReorderPolicyOperation operation) public List Generate(DropReorderPolicyOperation operation) { + string qualifiedTableName = sqlHelper.Regclass(operation.TableName, operation.Schema); + List statements = [ - $"SELECT remove_reorder_policy({sqlHelper.Regclass(operation.TableName)}, if_exists => true);" + $"SELECT remove_reorder_policy({qualifiedTableName}, if_exists => true);" ]; return statements; } @@ -125,18 +129,20 @@ private static List BuildAlterJobClauses(AlterReorderPolicyOperation ope return clauses; } - private static string BuildAlterJobSql(string tableName, IEnumerable clauses) + private static string BuildAlterJobSql(string tableName, string schema, IEnumerable clauses) { // Note: hypertable_name is a varchar column, so it compares against a string literal, not a regclass. return $@" SELECT alter_job(job_id, {string.Join(", ", clauses)}) FROM timescaledb_information.jobs - WHERE proc_name = 'policy_reorder' AND hypertable_name = '{tableName}';".Trim(); + WHERE proc_name = 'policy_reorder' AND hypertable_schema = '{schema}' AND hypertable_name = '{tableName}';".Trim(); } - private string BuildAddReorderPolicySql(string tableName, string indexName, DateTime? initialStart) + private string BuildAddReorderPolicySql(string tableName, string schema, string indexName, DateTime? initialStart) { - string baseSql = $"SELECT add_reorder_policy({sqlHelper.Regclass(tableName)}, '{indexName}'"; + string qualifiedTableName = sqlHelper.Regclass(tableName, schema); + + string baseSql = $"SELECT add_reorder_policy({qualifiedTableName}, '{indexName}'"; List optionalArgs = []; diff --git a/CmdScale.EntityFrameworkCore.TimescaleDB/Generators/SqlBuilderHelper.cs b/CmdScale.EntityFrameworkCore.TimescaleDB/Generators/SqlBuilderHelper.cs index 24369ee..8c25240 100644 --- a/CmdScale.EntityFrameworkCore.TimescaleDB/Generators/SqlBuilderHelper.cs +++ b/CmdScale.EntityFrameworkCore.TimescaleDB/Generators/SqlBuilderHelper.cs @@ -33,14 +33,14 @@ public static void BuildQueryString(List statements, IndentedStringBuild } } - public string Regclass(string tableName) + public string Regclass(string tableName, string schema = DefaultValues.DefaultSchema) { - return $"'{quoteString}{tableName}{quoteString}'"; + return $"'{schema}.{quoteString}{tableName}{quoteString}'"; } - public string QualifiedIdentifier(string tableName) + public string QualifiedIdentifier(string tableName, string schema = DefaultValues.DefaultSchema) { - return $"{quoteString}{tableName}{quoteString}"; + return $"{quoteString}{schema}{quoteString}.{quoteString}{tableName}{quoteString}"; } } } diff --git a/CmdScale.EntityFrameworkCore.TimescaleDB/Operations/AddReorderPolicyOperation.cs b/CmdScale.EntityFrameworkCore.TimescaleDB/Operations/AddReorderPolicyOperation.cs index 3fd646f..529cef8 100644 --- a/CmdScale.EntityFrameworkCore.TimescaleDB/Operations/AddReorderPolicyOperation.cs +++ b/CmdScale.EntityFrameworkCore.TimescaleDB/Operations/AddReorderPolicyOperation.cs @@ -5,6 +5,7 @@ namespace CmdScale.EntityFrameworkCore.TimescaleDB.Operations public class AddReorderPolicyOperation : MigrationOperation { public string TableName { get; set; } = string.Empty; + public string Schema { get; set; } = string.Empty; public string IndexName { get; set; } = string.Empty; public DateTime? InitialStart { get; set; } public string? ScheduleInterval { get; set; } diff --git a/CmdScale.EntityFrameworkCore.TimescaleDB/Operations/AlterHypertableOperation.cs b/CmdScale.EntityFrameworkCore.TimescaleDB/Operations/AlterHypertableOperation.cs index 2af4570..6f99e22 100644 --- a/CmdScale.EntityFrameworkCore.TimescaleDB/Operations/AlterHypertableOperation.cs +++ b/CmdScale.EntityFrameworkCore.TimescaleDB/Operations/AlterHypertableOperation.cs @@ -5,6 +5,7 @@ namespace CmdScale.EntityFrameworkCore.TimescaleDB.Operations public class AlterHypertableOperation : MigrationOperation { public string TableName { get; set; } = string.Empty; + public string Schema { get; set; } = string.Empty; public string ChunkTimeInterval { get; set; } = string.Empty; public bool EnableCompression { get; set; } diff --git a/CmdScale.EntityFrameworkCore.TimescaleDB/Operations/AlterReorderPolicyOperation.cs b/CmdScale.EntityFrameworkCore.TimescaleDB/Operations/AlterReorderPolicyOperation.cs index 436649b..46ba7d6 100644 --- a/CmdScale.EntityFrameworkCore.TimescaleDB/Operations/AlterReorderPolicyOperation.cs +++ b/CmdScale.EntityFrameworkCore.TimescaleDB/Operations/AlterReorderPolicyOperation.cs @@ -5,6 +5,7 @@ namespace CmdScale.EntityFrameworkCore.TimescaleDB.Operations public class AlterReorderPolicyOperation : MigrationOperation { public string TableName { get; set; } = string.Empty; + public string Schema { get; set; } = string.Empty; public string IndexName { get; set; } = string.Empty; public DateTime? InitialStart { get; set; } diff --git a/CmdScale.EntityFrameworkCore.TimescaleDB/Operations/CreateHypertableOperation.cs b/CmdScale.EntityFrameworkCore.TimescaleDB/Operations/CreateHypertableOperation.cs index 330b910..30850d4 100644 --- a/CmdScale.EntityFrameworkCore.TimescaleDB/Operations/CreateHypertableOperation.cs +++ b/CmdScale.EntityFrameworkCore.TimescaleDB/Operations/CreateHypertableOperation.cs @@ -6,6 +6,7 @@ namespace CmdScale.EntityFrameworkCore.TimescaleDB.Operations public class CreateHypertableOperation : MigrationOperation { public string TableName { get; set; } = string.Empty; + public string Schema { get; set; } = string.Empty; public string TimeColumnName { get; set; } = string.Empty; public string ChunkTimeInterval { get; set; } = string.Empty; public bool EnableCompression { get; set; } diff --git a/CmdScale.EntityFrameworkCore.TimescaleDB/Operations/DropReorderPolicyOperation.cs b/CmdScale.EntityFrameworkCore.TimescaleDB/Operations/DropReorderPolicyOperation.cs index e17a7a6..235797f 100644 --- a/CmdScale.EntityFrameworkCore.TimescaleDB/Operations/DropReorderPolicyOperation.cs +++ b/CmdScale.EntityFrameworkCore.TimescaleDB/Operations/DropReorderPolicyOperation.cs @@ -5,5 +5,6 @@ namespace CmdScale.EntityFrameworkCore.TimescaleDB.Operations public class DropReorderPolicyOperation : MigrationOperation { public string TableName { get; set; } = string.Empty; + public string Schema { get; set; } = string.Empty; } } diff --git a/CmdScale.EntityFrameworkCore.TimescaleDB/TimescaleMigrationsModelDiffer.cs b/CmdScale.EntityFrameworkCore.TimescaleDB/TimescaleMigrationsModelDiffer.cs index d686b88..2e4dcbb 100644 --- a/CmdScale.EntityFrameworkCore.TimescaleDB/TimescaleMigrationsModelDiffer.cs +++ b/CmdScale.EntityFrameworkCore.TimescaleDB/TimescaleMigrationsModelDiffer.cs @@ -37,13 +37,14 @@ public override IReadOnlyList GetDifferences(IRelationalMode List sourceHypertables = [.. GetHypertables(source)]; // Identify new hypertables - List newHypertables = [.. targetHypertables.Where(t => !sourceHypertables.Any(s => s.TableName == t.TableName))]; + List newHypertables = [.. targetHypertables.Where(t => !sourceHypertables.Any(s => s.TableName == t.TableName && s.Schema == t.Schema))]; foreach (CreateHypertableOperation? hypertable in newHypertables) { int createTableOpIndex = operations.FindIndex(op => op is CreateTableOperation createTable && - createTable.Name == hypertable.TableName); + createTable.Name == hypertable.TableName && + createTable.Schema == hypertable.Schema); if (createTableOpIndex != -1) { @@ -55,8 +56,8 @@ op is CreateTableOperation createTable && var updatedHypertables = targetHypertables .Join( sourceHypertables, - target => target.TableName, - source => source.TableName, + target => (target.Schema, target.TableName), + source => (source.Schema, source.TableName), (target, source) => new { Target = target, Source = source } ) .Where(x => @@ -71,6 +72,7 @@ op is CreateTableOperation createTable && AlterHypertableOperation alterOperation = new() { TableName = hypertable.Target.TableName, + Schema = hypertable.Target.Schema, ChunkTimeInterval = hypertable.Target.ChunkTimeInterval, EnableCompression = hypertable.Target.EnableCompression, ChunkSkipColumns = hypertable.Target.ChunkSkipColumns, @@ -88,15 +90,15 @@ op is CreateTableOperation createTable && List targetPolicies = [.. GetReorderPolicies(target)]; // Identiy new reorder policies - IEnumerable newReorderPolicies = targetPolicies.Where(t => !sourcePolicies.Any(s => s.TableName == t.TableName)); + IEnumerable newReorderPolicies = targetPolicies.Where(t => !sourcePolicies.Any(s => s.TableName == t.TableName && s.Schema == t.Schema)); operations.AddRange(newReorderPolicies); // Identify updated reorder policies var updatedReorderPolicies = targetPolicies .Join( sourcePolicies, - targetPolicy => targetPolicy.TableName, - sourcePolicy => sourcePolicy.TableName, + targetPolicy => (targetPolicy.Schema, targetPolicy.TableName), + sourcePolicy => (sourcePolicy.Schema, sourcePolicy.TableName), (targetPolicy, sourcePolicy) => new { Target = targetPolicy, Source = sourcePolicy } ) .Where(x => @@ -113,6 +115,7 @@ op is CreateTableOperation createTable && operations.Add(new AlterReorderPolicyOperation { TableName = policy.Target.TableName, + Schema = policy.Target.Schema, IndexName = policy.Target.IndexName, InitialStart = policy.Target.InitialStart, ScheduleInterval = policy.Target.ScheduleInterval, @@ -130,8 +133,8 @@ op is CreateTableOperation createTable && } IEnumerable removedReorderPolicies = sourcePolicies - .Where(s => !targetPolicies.Any(t => t.TableName == s.TableName)) - .Select(p => new DropReorderPolicyOperation { TableName = p.TableName }); + .Where(s => !targetPolicies.Any(t => t.TableName == s.TableName && t.Schema == s.Schema)) + .Select(p => new DropReorderPolicyOperation { TableName = p.TableName, Schema = p.Schema }); operations.AddRange(removedReorderPolicies); return operations; @@ -171,6 +174,7 @@ private static IEnumerable GetHypertables(IRelational yield return new CreateHypertableOperation { TableName = entityType.GetTableName()!, + Schema = entityType.GetSchema() ?? DefaultValues.DefaultSchema, TimeColumnName = timeColumnName, ChunkTimeInterval = chunkTimeInterval ?? DefaultValues.ChunkTimeInterval, EnableCompression = enableCompression, @@ -199,6 +203,7 @@ private static IEnumerable GetReorderPolicies(IRelati yield return new AddReorderPolicyOperation { TableName = entityType.GetTableName()!, + Schema = entityType.GetSchema() ?? DefaultValues.DefaultSchema, IndexName = indexName!, InitialStart = initialStart, ScheduleInterval = entityType.FindAnnotation(ReorderPolicyAnnotations.ScheduleInterval)?.Value as string ?? DefaultValues.ReorderPolicyScheduleInterval,