-
Notifications
You must be signed in to change notification settings - Fork 4
Expand file tree
/
Copy pathSchemaController.cs
More file actions
129 lines (112 loc) · 5.36 KB
/
Copy pathSchemaController.cs
File metadata and controls
129 lines (112 loc) · 5.36 KB
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
using k8s.Models;
using KubeOps.Abstractions.Rbac;
using KubeOps.Abstractions.Reconciliation;
using KubeOps.Abstractions.Reconciliation.Controller;
using KubeOps.KubernetesClient;
using Microsoft.Data.SqlClient;
using SqlServerOperator.Controllers.Services;
using SqlServerOperator.Entities.V1Alpha1;
using System.Text;
namespace SqlServerOperator.Controllers.V1Alpha1;
[EntityRbac(typeof(V1Alpha1SQLServerSchema), Verbs = RbacVerb.All)]
public class SQLServerSchemaController(
ILogger<SQLServerSchemaController> logger,
IKubernetesClient kubernetesClient,
ISqlServerEndpointService sqlServerEndpointService,
ISqlExecutor sqlExecutor
) : IEntityController<V1Alpha1SQLServerSchema>
{
public async Task<ReconciliationResult<V1Alpha1SQLServerSchema>> ReconcileAsync(V1Alpha1SQLServerSchema entity, CancellationToken cancellationToken)
{
logger.LogInformation("Reconciling SQLServerSchema: {Name}", entity.Metadata.Name);
try
{
// Try ExternalSQLServer first
var externalServer = await kubernetesClient.GetAsync<V1Alpha1ExternalSQLServer>(entity.Spec.InstanceName, entity.Metadata.NamespaceProperty);
string secretName;
if (externalServer is not null)
{
secretName = externalServer.Spec.SecretName;
}
else
{
// Fall back to internal SQLServer
var sqlServer = await kubernetesClient.GetAsync<V1Alpha1SQLServer>(entity.Spec.InstanceName, entity.Metadata.NamespaceProperty);
if (sqlServer is null)
{
throw new Exception($"SQLServer or ExternalSQLServer instance '{entity.Spec.InstanceName}' not found.");
}
secretName = sqlServer.Spec.SecretName ?? $"{sqlServer.Metadata.Name}-secret";
}
var server = await sqlServerEndpointService.GetSqlServerEndpointAsync(entity.Spec.InstanceName, entity.Metadata.NamespaceProperty);
var (username, password) = await GetSqlServerCredentialsAsync(secretName, entity.Metadata.NamespaceProperty);
await EnsureSchemaExistsAsync(
entity.Spec.DatabaseName,
entity.Spec.SchemaName,
entity.Spec.SchemaOwner,
server,
username,
password);
entity.Status ??= new();
entity.Status.State = "Ready";
entity.Status.Message = "Schema ensured.";
entity.Status.LastChecked = DateTime.UtcNow;
await kubernetesClient.UpdateStatusAsync(entity);
return ReconciliationResult<V1Alpha1SQLServerSchema>.Success(entity, TimeSpan.FromMinutes(5));
}
catch (Exception ex)
{
logger.LogError(ex, "Error during reconciliation of SQLServerSchema: {Name}", entity.Metadata.Name);
entity.Status ??= new();
entity.Status.State = "Error";
entity.Status.Message = ex.Message;
entity.Status.LastChecked = DateTime.UtcNow;
await kubernetesClient.UpdateStatusAsync(entity);
return ReconciliationResult<V1Alpha1SQLServerSchema>.Failure(entity, ex.Message, ex, TimeSpan.FromMinutes(1));
}
}
public Task<ReconciliationResult<V1Alpha1SQLServerSchema>> DeletedAsync(V1Alpha1SQLServerSchema entity, CancellationToken cancellationToken)
{
logger.LogInformation("Deleted SQLServerSchema: {Name}", entity.Metadata.Name);
return Task.FromResult(ReconciliationResult<V1Alpha1SQLServerSchema>.Success(entity));
}
private async Task<(string username, string password)> GetSqlServerCredentialsAsync(string secretName, string namespaceName)
{
var secret = await kubernetesClient.GetAsync<V1Secret>(secretName, namespaceName);
if (secret?.Data is null || !secret.Data.ContainsKey("password"))
{
throw new Exception($"Secret '{secretName}' does not contain the expected 'password' key.");
}
var password = Encoding.UTF8.GetString(secret.Data["password"]);
var username = "sa";
return (username, password);
}
private async Task EnsureSchemaExistsAsync(string databaseName, string schemaName, string schemaOwner, string server, string username, string password)
{
var builder = new SqlConnectionStringBuilder
{
DataSource = server,
UserID = username,
Password = password,
InitialCatalog = databaseName,
TrustServerCertificate = true,
Encrypt = false,
};
var schemaExistsCommandText = @"
IF NOT EXISTS (
SELECT schema_name
FROM information_schema.schemata
WHERE schema_name = @SchemaName
)
BEGIN
DECLARE @sql NVARCHAR(MAX) = N'CREATE SCHEMA [' + @SchemaName + '] AUTHORIZATION [' + @SchemaOwner + ']';
EXEC sp_executesql @sql;
END";
var parameters = new Dictionary<string, object>
{
["@SchemaName"] = schemaName,
["@SchemaOwner"] = schemaOwner
};
await sqlExecutor.ExecuteNonQueryAsync(builder.ConnectionString, schemaExistsCommandText, parameters);
}
}