using Microsoft.Data.Sqlite; using UnrealDemoScanner.Web.Infrastructure; using UnrealDemoScanner.Web.Models; namespace UnrealDemoScanner.Web.Services; public sealed class SqliteScanRepository : IScanRepository { private readonly AppPaths _paths; public SqliteScanRepository(AppPaths paths) { _paths = paths; } public async Task InitializeAsync(CancellationToken cancellationToken = default) { Directory.CreateDirectory(_paths.DataRoot); Directory.CreateDirectory(_paths.UploadsRoot); Directory.CreateDirectory(_paths.ResultsRoot); await using var connection = OpenConnection(); await connection.OpenAsync(cancellationToken); var sql = """ CREATE TABLE IF NOT EXISTS jobs ( id TEXT PRIMARY KEY, original_file_name TEXT NOT NULL, stored_file_path TEXT NOT NULL, status TEXT NOT NULL, submitted_at_utc TEXT NOT NULL, started_at_utc TEXT NULL, completed_at_utc TEXT NULL, exit_code INTEGER NULL, error_message TEXT NULL, output_log_path TEXT NULL, detected_count INTEGER NOT NULL DEFAULT 0, warning_count INTEGER NOT NULL DEFAULT 0, info_count INTEGER NOT NULL DEFAULT 0, duration_ms INTEGER NULL ); CREATE TABLE IF NOT EXISTS findings ( id INTEGER PRIMARY KEY AUTOINCREMENT, job_id TEXT NOT NULL, kind TEXT NOT NULL, line_number INTEGER NOT NULL, message TEXT NOT NULL, created_at_utc TEXT NOT NULL, FOREIGN KEY(job_id) REFERENCES jobs(id) ON DELETE CASCADE ); CREATE INDEX IF NOT EXISTS idx_jobs_submitted_at ON jobs(submitted_at_utc DESC); CREATE INDEX IF NOT EXISTS idx_findings_job_id ON findings(job_id); """; await using var command = connection.CreateCommand(); command.CommandText = sql; await command.ExecuteNonQueryAsync(cancellationToken); } public async Task CreateJobAsync(ScanJob job, CancellationToken cancellationToken = default) { const string sql = """ INSERT INTO jobs ( id, original_file_name, stored_file_path, status, submitted_at_utc ) VALUES ( $id, $original_file_name, $stored_file_path, $status, $submitted_at_utc ); """; await using var connection = OpenConnection(); await connection.OpenAsync(cancellationToken); await using var command = connection.CreateCommand(); command.CommandText = sql; command.Parameters.AddWithValue("$id", job.Id); command.Parameters.AddWithValue("$original_file_name", job.OriginalFileName); command.Parameters.AddWithValue("$stored_file_path", job.StoredFilePath); command.Parameters.AddWithValue("$status", job.Status.ToString()); command.Parameters.AddWithValue("$submitted_at_utc", job.SubmittedAtUtc.ToString("O")); await command.ExecuteNonQueryAsync(cancellationToken); } public async Task GetJobAsync(string id, CancellationToken cancellationToken = default) { const string sql = "SELECT * FROM jobs WHERE id = $id LIMIT 1;"; await using var connection = OpenConnection(); await connection.OpenAsync(cancellationToken); await using var command = connection.CreateCommand(); command.CommandText = sql; command.Parameters.AddWithValue("$id", id); await using var reader = await command.ExecuteReaderAsync(cancellationToken); if (!await reader.ReadAsync(cancellationToken)) { return null; } return MapJob(reader); } public async Task> ListJobsAsync(int limit, CancellationToken cancellationToken = default) { const string sql = """ SELECT * FROM jobs ORDER BY submitted_at_utc DESC LIMIT $limit; """; var jobs = new List(); await using var connection = OpenConnection(); await connection.OpenAsync(cancellationToken); await using var command = connection.CreateCommand(); command.CommandText = sql; command.Parameters.AddWithValue("$limit", limit); await using var reader = await command.ExecuteReaderAsync(cancellationToken); while (await reader.ReadAsync(cancellationToken)) { jobs.Add(MapJob(reader)); } return jobs; } public async Task> GetFindingsAsync(string jobId, CancellationToken cancellationToken = default) { const string sql = """ SELECT id, job_id, kind, line_number, message, created_at_utc FROM findings WHERE job_id = $job_id ORDER BY id ASC; """; var findings = new List(); await using var connection = OpenConnection(); await connection.OpenAsync(cancellationToken); await using var command = connection.CreateCommand(); command.CommandText = sql; command.Parameters.AddWithValue("$job_id", jobId); await using var reader = await command.ExecuteReaderAsync(cancellationToken); while (await reader.ReadAsync(cancellationToken)) { findings.Add(new ScanFinding { Id = reader.GetInt64(0), JobId = reader.GetString(1), Kind = reader.GetString(2), LineNumber = reader.GetInt32(3), Message = reader.GetString(4), CreatedAtUtc = DateTime.Parse(reader.GetString(5)).ToUniversalTime() }); } return findings; } public async Task MarkRunningAsync(string id, DateTime startedAtUtc, CancellationToken cancellationToken = default) { const string sql = """ UPDATE jobs SET status = $status, started_at_utc = $started_at_utc WHERE id = $id; """; await using var connection = OpenConnection(); await connection.OpenAsync(cancellationToken); await using var command = connection.CreateCommand(); command.CommandText = sql; command.Parameters.AddWithValue("$status", ScanJobStatus.Running.ToString()); command.Parameters.AddWithValue("$started_at_utc", startedAtUtc.ToString("O")); command.Parameters.AddWithValue("$id", id); await command.ExecuteNonQueryAsync(cancellationToken); } public async Task CompleteAsync( string id, ScanJobStatus status, int? exitCode, string? errorMessage, string? outputLogPath, long durationMs, int detectedCount, int warningCount, int infoCount, IReadOnlyList findings, CancellationToken cancellationToken = default) { await using var connection = OpenConnection(); await connection.OpenAsync(cancellationToken); await using var transaction = (SqliteTransaction)await connection.BeginTransactionAsync(cancellationToken); var deleteFindings = connection.CreateCommand(); deleteFindings.Transaction = transaction; deleteFindings.CommandText = "DELETE FROM findings WHERE job_id = $job_id;"; deleteFindings.Parameters.AddWithValue("$job_id", id); await deleteFindings.ExecuteNonQueryAsync(cancellationToken); foreach (var finding in findings) { var insertFinding = connection.CreateCommand(); insertFinding.Transaction = transaction; insertFinding.CommandText = """ INSERT INTO findings (job_id, kind, line_number, message, created_at_utc) VALUES ($job_id, $kind, $line_number, $message, $created_at_utc); """; insertFinding.Parameters.AddWithValue("$job_id", id); insertFinding.Parameters.AddWithValue("$kind", finding.Kind); insertFinding.Parameters.AddWithValue("$line_number", finding.LineNumber); insertFinding.Parameters.AddWithValue("$message", finding.Message); insertFinding.Parameters.AddWithValue("$created_at_utc", finding.CreatedAtUtc.ToString("O")); await insertFinding.ExecuteNonQueryAsync(cancellationToken); } var completeCommand = connection.CreateCommand(); completeCommand.Transaction = transaction; completeCommand.CommandText = """ UPDATE jobs SET status = $status, completed_at_utc = $completed_at_utc, exit_code = $exit_code, error_message = $error_message, output_log_path = $output_log_path, duration_ms = $duration_ms, detected_count = $detected_count, warning_count = $warning_count, info_count = $info_count WHERE id = $id; """; completeCommand.Parameters.AddWithValue("$status", status.ToString()); completeCommand.Parameters.AddWithValue("$completed_at_utc", DateTime.UtcNow.ToString("O")); completeCommand.Parameters.AddWithValue("$exit_code", (object?)exitCode ?? DBNull.Value); completeCommand.Parameters.AddWithValue("$error_message", (object?)errorMessage ?? DBNull.Value); completeCommand.Parameters.AddWithValue("$output_log_path", (object?)outputLogPath ?? DBNull.Value); completeCommand.Parameters.AddWithValue("$duration_ms", durationMs); completeCommand.Parameters.AddWithValue("$detected_count", detectedCount); completeCommand.Parameters.AddWithValue("$warning_count", warningCount); completeCommand.Parameters.AddWithValue("$info_count", infoCount); completeCommand.Parameters.AddWithValue("$id", id); await completeCommand.ExecuteNonQueryAsync(cancellationToken); await transaction.CommitAsync(cancellationToken); } private SqliteConnection OpenConnection() { var connectionString = new SqliteConnectionStringBuilder { DataSource = _paths.DatabasePath, ForeignKeys = true }.ToString(); return new SqliteConnection(connectionString); } private static ScanJob MapJob(SqliteDataReader reader) { return new ScanJob { Id = reader.GetString(reader.GetOrdinal("id")), OriginalFileName = reader.GetString(reader.GetOrdinal("original_file_name")), StoredFilePath = reader.GetString(reader.GetOrdinal("stored_file_path")), Status = Enum.Parse(reader.GetString(reader.GetOrdinal("status"))), SubmittedAtUtc = DateTime.Parse(reader.GetString(reader.GetOrdinal("submitted_at_utc"))).ToUniversalTime(), StartedAtUtc = ReadNullableDateTime(reader, "started_at_utc"), CompletedAtUtc = ReadNullableDateTime(reader, "completed_at_utc"), ExitCode = ReadNullableInt(reader, "exit_code"), ErrorMessage = ReadNullableString(reader, "error_message"), OutputLogPath = ReadNullableString(reader, "output_log_path"), DetectedCount = reader.GetInt32(reader.GetOrdinal("detected_count")), WarningCount = reader.GetInt32(reader.GetOrdinal("warning_count")), InfoCount = reader.GetInt32(reader.GetOrdinal("info_count")), DurationMs = ReadNullableLong(reader, "duration_ms") }; } private static DateTime? ReadNullableDateTime(SqliteDataReader reader, string column) { var ordinal = reader.GetOrdinal(column); if (reader.IsDBNull(ordinal)) { return null; } return DateTime.Parse(reader.GetString(ordinal)).ToUniversalTime(); } private static string? ReadNullableString(SqliteDataReader reader, string column) { var ordinal = reader.GetOrdinal(column); return reader.IsDBNull(ordinal) ? null : reader.GetString(ordinal); } private static int? ReadNullableInt(SqliteDataReader reader, string column) { var ordinal = reader.GetOrdinal(column); return reader.IsDBNull(ordinal) ? null : reader.GetInt32(ordinal); } private static long? ReadNullableLong(SqliteDataReader reader, string column) { var ordinal = reader.GetOrdinal(column); return reader.IsDBNull(ordinal) ? null : reader.GetInt64(ordinal); } }