tbf/TBF/Rig/Output/DataStorage/UniDataStorageWriter/Writers/DatabaseWriter .cs

145 lines
4.5 KiB
C#

using System;
using System.Data.SqlClient;
using System.Linq;
using TBF.Rig.Output.DataStorage.UniDataStorageWriter.Interfaces;
namespace TBF.Rig.Output.DataStorage.UniDataStorageWriter.Writers
{
public class DatabaseWriter : IDataStorageWriter
{
private readonly WriterCfg cfg;
public DatabaseWriter(WriterCfg cfg)
{
this.cfg = cfg ?? throw new ArgumentNullException(nameof(cfg));
}
public WriterDiagnosticResult TestSource(bool validateOnly)
{
try
{
using (var conn = new SqlConnection(cfg.DataSource))
{
conn.Open();
if (!validateOnly)
{
using (var cmd = new SqlCommand("SELECT 1", conn))
cmd.ExecuteScalar();
}
}
return Ok("Database connection successful.");
}
catch (Exception ex)
{
return Fail("Database connection failed: " + ex.Message);
}
}
public WriterDiagnosticResult WriteData(DataWriteRequest request)
{
if (request == null)
throw new ArgumentNullException(nameof(request));
if (string.IsNullOrWhiteSpace(request.TargetName))
return Fail("TargetName (table) must be defined.");
switch (request.Mode)
{
case WriteMode.Insert:
return ExecuteInsert(request);
case WriteMode.Update:
return ExecuteUpdate(request);
default:
return Fail("Mode not supported: " + request.Mode);
}
}
private WriterDiagnosticResult ExecuteInsert(DataWriteRequest request)
{
if (request.InsertItems.Count == 0)
return Fail("No insert items.");
string[] columns = request.InsertItems.Select(i => i.ColumnName).ToArray();
string[] values = request.InsertItems.Select(i => ToSql(i.Value)).ToArray();
string sql = $"INSERT INTO {request.TargetName} ({string.Join(", ", columns)}) VALUES ({string.Join(", ", values)})";
return ExecuteSql(sql);
}
private WriterDiagnosticResult ExecuteUpdate(DataWriteRequest request)
{
if (request.UpdateItems.Count == 0)
return Fail("No update items.");
int total = 0;
string executed = "";
using (var conn = new SqlConnection(cfg.DataSource))
{
conn.Open();
foreach (var item in request.UpdateItems)
{
string sql = $"UPDATE {request.TargetName} " +
$"SET {item.SetParameterName} = {ToSql(item.SetValue)} " +
$"WHERE {item.WhereParameterName} = {ToSql(item.WhereValue)}";
using (var cmd = new SqlCommand(sql, conn))
{
total += cmd.ExecuteNonQuery();
}
executed += sql + Environment.NewLine;
}
}
return new WriterDiagnosticResult
{
Success = true,
Message = "Update OK. Rows: " + total,
ExecutedTemplate = executed
};
}
private WriterDiagnosticResult ExecuteSql(string sql)
{
using (var conn = new SqlConnection(cfg.DataSource))
{
conn.Open();
using (var cmd = new SqlCommand(sql, conn))
{
int rows = cmd.ExecuteNonQuery();
return new WriterDiagnosticResult
{
Success = true,
Message = "Insert OK. Rows: " + rows,
ExecutedTemplate = sql
};
}
}
}
private string ToSql(string val)
{
if (val == null) return "NULL";
return "'" + val.Replace("'", "''") + "'";
}
private WriterDiagnosticResult Ok(string msg)
{
return new WriterDiagnosticResult { Success = true, Message = msg };
}
private WriterDiagnosticResult Fail(string msg)
{
return new WriterDiagnosticResult { Success = false, Message = msg };
}
}
}