using System.ComponentModel.DataAnnotations.Schema; using System.Reflection; using Npgsql; using StickyBoard.Core.Common; using StickyBoard.Core.Models.Base; namespace StickyBoard.Core.Repositories.Base; public abstract class RepositoryBase : IRepository, ISyncRepository where T : class, IEntity, new() { private readonly NpgsqlDataSource _db; protected RepositoryBase(NpgsqlDataSource db) { _db = db; } protected string Table => typeof(T).GetCustomAttribute()?.Name ?? typeof(T).Name.ToLowerInvariant(); // --------------------------------------------------------------------- // Soft Delete Helpers // --------------------------------------------------------------------- private bool SoftEnabled => typeof(ISoftDeletable).IsAssignableFrom(typeof(T)); private bool IncludeDeleted => this is IAllowDeleted allow && allow.IncludeDeleted; // --------------------------------------------------------------------- // GET BY ID // --------------------------------------------------------------------- public async Task GetByIdAsync(Guid id, CancellationToken ct) { var sql = ApplySoftDeleteFilter($"SELECT * FROM {Table} WHERE id = @id"); await using var c = await Conn(ct); await using var cmd = new NpgsqlCommand(sql, c); cmd.Parameters.AddWithValue("id", id); await using var r = await cmd.ExecuteReaderAsync(ct); return await r.ReadAsync(ct) ? MapRow(r) : null; } public async Task GetByIdIncludingDeletedAsync(Guid id, CancellationToken ct) { var sql = $"SELECT * FROM {Table} WHERE id = @id"; await using var c = await Conn(ct); await using var cmd = new NpgsqlCommand(sql, c); cmd.Parameters.AddWithValue("id", id); await using var r = await cmd.ExecuteReaderAsync(ct); return await r.ReadAsync(ct) ? MapRow(r) : null; } // --------------------------------------------------------------------- // GET ALL // --------------------------------------------------------------------- public async Task> GetAllAsync(CancellationToken ct) { var sql = ApplySoftDeleteFilter($"SELECT * FROM {Table}"); await using var c = await Conn(ct); await using var cmd = new NpgsqlCommand(sql, c); await using var r = await cmd.ExecuteReaderAsync(ct); return await MapListAsync(r, ct); } // --------------------------------------------------------------------- // EXISTS // --------------------------------------------------------------------- public async Task ExistsAsync(Guid id, CancellationToken ct) { var sql = ApplySoftDeleteFilter($"SELECT 1 FROM {Table} WHERE id = @id"); await using var c = await Conn(ct); await using var cmd = new NpgsqlCommand(sql, c); cmd.Parameters.AddWithValue("id", id); return await cmd.ExecuteScalarAsync(ct) is not null; } // --------------------------------------------------------------------- // COUNT // --------------------------------------------------------------------- public async Task CountAsync(CancellationToken ct) { var sql = ApplySoftDeleteFilter($"SELECT COUNT(*) FROM {Table}"); await using var c = await Conn(ct); await using var cmd = new NpgsqlCommand(sql, c); return Convert.ToInt32(await cmd.ExecuteScalarAsync(ct)); } // --------------------------------------------------------------------- // PAGING (single-query with window function) // --------------------------------------------------------------------- public async Task> GetPagedAsync(int limit, int offset, CancellationToken ct) { var sql = ApplySoftDeleteFilter($@" SELECT *, COUNT(*) OVER() AS total_count FROM {Table} ORDER BY updated_at DESC LIMIT @limit OFFSET @offset"); await using var c = await Conn(ct); await using var cmd = new NpgsqlCommand(sql, c); cmd.Parameters.AddWithValue("limit", limit); cmd.Parameters.AddWithValue("offset", offset); await using var r = await cmd.ExecuteReaderAsync(ct); var items = new List(); var total = 0; while (await r.ReadAsync(ct)) { items.Add(MapRow(r)); total = r.GetInt32(r.GetOrdinal("total_count")); } return PagedResult.Create(items, total, limit, offset); } // --------------------------------------------------------------------- // CREATE (implemented by subclasses) // --------------------------------------------------------------------- public abstract Task CreateAsync(T entity, CancellationToken ct); // --------------------------------------------------------------------- // UPDATE (must enforce concurrency if entity is versioned) // --------------------------------------------------------------------- public abstract Task UpdateAsync(T entity, CancellationToken ct); // --------------------------------------------------------------------- // DELETE // --------------------------------------------------------------------- public async Task DeleteAsync(Guid id, CancellationToken ct) { await using var c = await Conn(ct); if (SoftEnabled) { var sql = $"UPDATE {Table} SET deleted_at = NOW() WHERE id = @id"; await using var cmd = new NpgsqlCommand(sql, c); cmd.Parameters.AddWithValue("id", id); return await cmd.ExecuteNonQueryAsync(ct) > 0; } var hard = $"DELETE FROM {Table} WHERE id = @id"; await using var cmdHard = new NpgsqlCommand(hard, c); cmdHard.Parameters.AddWithValue("id", id); return await cmdHard.ExecuteNonQueryAsync(ct) > 0; } // --------------------------------------------------------------------- // SYNC (delta + single-query paging) // --------------------------------------------------------------------- public async Task> GetUpdatedSinceAsync(DateTime since, CancellationToken ct) { var sql = ApplySoftDeleteFilter($@" SELECT * FROM {Table} WHERE updated_at > @since"); await using var c = await Conn(ct); await using var cmd = new NpgsqlCommand(sql, c); cmd.Parameters.AddWithValue("since", since); await using var r = await cmd.ExecuteReaderAsync(ct); return await MapListAsync(r, ct); } public async Task> GetUpdatedSincePagedAsync( DateTime since, int limit, int offset, CancellationToken ct) { var sql = ApplySoftDeleteFilter($@" SELECT *, COUNT(*) OVER() AS total_count FROM {Table} WHERE updated_at > @since ORDER BY updated_at DESC LIMIT @limit OFFSET @offset"); await using var c = await Conn(ct); await using var cmd = new NpgsqlCommand(sql, c); cmd.Parameters.AddWithValue("since", since); cmd.Parameters.AddWithValue("limit", limit); cmd.Parameters.AddWithValue("offset", offset); await using var r = await cmd.ExecuteReaderAsync(ct); var items = new List(); var total = 0; while (await r.ReadAsync(ct)) { items.Add(MapRow(r)); total = r.GetInt32(r.GetOrdinal("total_count")); } return PagedResult.Create(items, total, limit, offset); } protected ValueTask Conn(CancellationToken ct) { return _db.OpenConnectionAsync(ct); } protected string ApplySoftDeleteFilter(string sql) { if (!SoftEnabled || IncludeDeleted) return sql; return sql.Contains("WHERE", StringComparison.OrdinalIgnoreCase) ? sql + " AND deleted_at IS NULL" : sql + " WHERE deleted_at IS NULL"; } // --------------------------------------------------------------------- // Mapping // --------------------------------------------------------------------- protected virtual T MapRow(NpgsqlDataReader r) { return MappingHelper.MapEntity(r); } protected async Task> MapListAsync(NpgsqlDataReader r, CancellationToken ct) { var list = new List(); while (await r.ReadAsync(ct)) list.Add(MapRow(r)); return list; } protected string ConcurrencyWhere(T entity) { if (entity is IVersionedEntity v) return "id = @id AND version = @version"; return "id = @id"; } protected void BindConcurrencyParameters(NpgsqlCommand cmd, T entity) { if (entity is IVersionedEntity v) cmd.Parameters.AddWithValue("version", v.Version); } }