using System.Collections.Concurrent; using Npgsql; using NpgsqlTypes; namespace w4c_workflows.Services; /// /// Factory for creating instances scoped to a /// specific tenant. The caller provides the tenant ID explicitly (from auth /// middleware); the factory resolves directories and returns a source bound /// to that tenant's workflow directory. /// /// Two modes (automatic, config-driven): /// /// Forgejo-backed (preferred): when Forgejo:AdminToken is /// configured, each tenant's workflow files live in a private Forgejo repo /// (workflows-{tenantId}). The factory clones/pulls via /// and returns a /// backed by the local clone. /// This gives tenants real git history, branching, and remote sync. /// Filesystem-only (fallback): when Forgejo is not configured, /// uses plain directory copies seeded from the shared template dir. /// Preserves the original behavior for dev/self-hosted setups. /// /// /// Tenant isolation is mandatory: WorkflowSource:CopiesRoot must /// be set in configuration. The factory throws at construction time if it is /// missing, preventing any request from serving cross-tenant workflows. /// /// Registered as a Singleton (shared config, no per-request state); /// the created sources are lightweight and backed by filesystem operations. /// public sealed class WorkflowSourceFactory { private readonly ILoggerFactory _loggers; private readonly string _copiesRoot; private readonly ForgejoWorkflowRepoService? _forgejo; private readonly NpgsqlDataSource? _ds; private readonly GitRunner _git; private readonly TimeSpan _pullInterval; // Per-(tenant, repo) cache of the resolved source + Forgejo login so clone/pull // and the login lookup do not run on every authenticated HTTP request. A // per-key gate collapses concurrent refreshes into a single clone/pull. private readonly ConcurrentDictionary _cache = new(StringComparer.Ordinal); private readonly ConcurrentDictionary _gates = new(StringComparer.Ordinal); private sealed record CachedSource(IWorkflowSource Source, string? Login, DateTime CreatedAt); public WorkflowSourceFactory(IConfiguration config, IWebHostEnvironment env, ILoggerFactory loggers, GitRunner git, IHttpClientFactory httpFactory, NpgsqlDataSource? ds = null) { _loggers = loggers; _ds = ds; _git = git; _pullInterval = TimeSpan.FromSeconds( int.TryParse(config["WorkflowSource:PullIntervalSeconds"], out var seconds) && seconds > 0 ? seconds : 30); _copiesRoot = ResolveCopiesRoot( config["WorkflowSource:CopiesRoot"] ?? throw new InvalidOperationException( "WorkflowSource:CopiesRoot is not configured. " + "Per-tenant filesystem isolation is required — set it to a writable " + "directory path (e.g. \"/data/workflow-tenants\" or \".data/workflow-tenants\")."), env); // Create the Forgejo service when admin provisioning is configured. var forgejoLogger = loggers.CreateLogger(); var forgejoSvc = new ForgejoWorkflowRepoService(config, env, forgejoLogger, httpFactory); _forgejo = forgejoSvc.IsConfigured ? forgejoSvc : null; if (_forgejo != null) _loggers.CreateLogger() .LogInformation("Workflow source: Forgejo-backed mode (owner={Owner})", config["Forgejo:Owner"] ?? config["Forgejo:WorkflowRepoOwner"]); else _loggers.CreateLogger() .LogInformation("Workflow source: filesystem-only mode (no Forgejo admin token)"); } /// True when Forgejo-backed mode is active. public bool IsForgejoBacked => _forgejo != null; /// The Forgejo repo service (null when in filesystem-only mode). public ForgejoWorkflowRepoService? Forgejo => _forgejo; /// /// Creates a per-tenant workflow source. /// /// In Forgejo-backed mode: ensures the local clone exists (creating the /// Forgejo repo if needed), then returns a source backed by the clone. /// /// In filesystem-only mode: the tenant directory /// {CopiesRoot}/{sanitizedTenantId}/ is created empty on first access /// (never seeded from shared templates — workflows are strictly per-tenant). /// public IWorkflowSource Create(string tenantId) { var tenantDir = TenantSourceDir(tenantId); var logger = _loggers.CreateLogger(); return new PerTenantWorkflowSource(tenantId, tenantDir, logger, _git); } /// /// Resolves the absolute on-disk directory that holds a tenant's workflow /// files. In Forgejo-backed mode this is the per-user clone dir; in /// filesystem-only mode it is {CopiesRoot}/{sanitizedTenantId}/. /// Used by the workflow-file CRUD controller so edits always land in the same /// place reads from. /// public string ResolveTenantSourceDir(string tenantId, string? forgejoLogin = null, string? repoName = null) { if (_forgejo != null && !string.IsNullOrWhiteSpace(forgejoLogin)) return _forgejo.TenantCloneDir(tenantId, forgejoLogin, repoName); return TenantSourceDir(tenantId); } /// The per-tenant filesystem source directory (not Forgejo-backed). public string TenantSourceDir(string tenantId) => Path.GetFullPath(Path.Combine(_copiesRoot, Sanitize(tenantId))); /// /// Creates a per-tenant workflow source asynchronously. In Forgejo-backed /// mode this ensures the local clone exists (cloning from Forgejo if needed). /// In filesystem-only mode, delegates to . /// public async Task CreateAsync(string tenantId, string? repoName = null, CancellationToken ct = default) => (await ResolveAsync(tenantId, repoName, ct)).Source; /// /// Resolves the per-tenant source and its Forgejo login, caching the result /// for WorkflowSource:PullIntervalSeconds (default 30s) so clone/pull and /// the login lookup run at most once per interval instead of on every request. /// Concurrent refreshes for the same tenant+repo are collapsed into one. /// public async Task<(IWorkflowSource Source, string? Login)> ResolveAsync( string tenantId, string? repoName = null, CancellationToken ct = default) { if (_forgejo == null) return (Create(tenantId), null); var key = tenantId + "|" + (repoName ?? string.Empty); if (TryGetFresh(key, out var cached)) return (cached!.Source, cached.Login); var gate = _gates.GetOrAdd(key, _ => new SemaphoreSlim(1, 1)); await gate.WaitAsync(ct); try { if (TryGetFresh(key, out cached)) return (cached!.Source, cached.Login); // Resolve the user's Forgejo login from the tenant ID (forgejo_id). // The repo lives under the user's own account: {login}/{workflowRepoName}. var login = await ResolveForgejoLoginAsync(tenantId, ct); if (string.IsNullOrEmpty(login)) { _loggers.CreateLogger() .LogWarning("Could not resolve Forgejo login for tenant {TenantId}, falling back to filesystem", tenantId); return (Create(tenantId), null); } // Ensure the Forgejo repo exists and is cloned locally. var cloneDir = await _forgejo.EnsureCloneAsync(tenantId, login, repoName, ct); // Pull latest changes before reading (from the SELECTED repo, not the default). await _forgejo.PullAsync(tenantId, login, repoName, ct); var logger = _loggers.CreateLogger(); var source = new PerTenantWorkflowSource(tenantId, cloneDir, logger, _git); _cache[key] = new CachedSource(source, login, DateTime.UtcNow); return (source, login); } finally { gate.Release(); } } /// Drops the cached source for a tenant+repo, forcing a fresh clone/pull next time. public void Invalidate(string tenantId, string? repoName = null) => _cache.TryRemove(tenantId + "|" + (repoName ?? string.Empty), out _); private bool TryGetFresh(string key, out CachedSource? cached) { if (_cache.TryGetValue(key, out cached) && DateTime.UtcNow - cached.CreatedAt < _pullInterval) return true; cached = null; return false; } /// /// Resolves the Forgejo login for a tenant (forgejo_id) from the auth_users table. /// public async Task ResolveForgejoLoginAsync(string tenantId, CancellationToken ct) { if (_ds == null || !long.TryParse(tenantId, out var forgejoId)) return null; try { await using var conn = await _ds.OpenConnectionAsync(ct); await using var cmd = conn.CreateCommand(); cmd.CommandTimeout = 3; cmd.CommandText = "SELECT login FROM auth_users WHERE forgejo_id = @fid LIMIT 1"; cmd.Parameters.Add(new NpgsqlParameter("@fid", NpgsqlDbType.Bigint) { Value = forgejoId }); var result = await cmd.ExecuteScalarAsync(ct); return result as string; } catch (Exception ex) { _loggers.CreateLogger() .LogWarning(ex, "Could not resolve Forgejo login for tenant {TenantId}", tenantId); return null; } } /// /// Resolves WorkflowSource:CopiesRoot to an absolute path independent of /// the process CWD. A relative value (e.g. ../source-copies) is anchored /// to the app's content root instead of , /// so a worker started from a container/different working directory still lands on /// the same wizard root as the dev run. An absolute value is normalized as-is. /// internal static string ResolveCopiesRoot(string raw, IWebHostEnvironment? env) { if (Path.IsPathRooted(raw)) return Path.GetFullPath(raw); var basePath = env?.ContentRootPath; if (string.IsNullOrWhiteSpace(basePath)) basePath = Directory.GetCurrentDirectory(); return Path.GetFullPath(raw, basePath); } private static string Sanitize(string id) { if (string.IsNullOrWhiteSpace(id)) return "_"; var sb = new System.Text.StringBuilder(id.Length); foreach (var ch in id) sb.Append(char.IsLetterOrDigit(ch) || ch == '-' || ch == '_' ? ch : '_'); var result = sb.ToString(); return result.Length > 120 ? result[..120] : result; } } /// /// Per-tenant workflow YAML source with directory isolation. Each tenant gets /// its own directory at {CopiesRoot}/{tenantId}/, initialized from the /// shared template directory on first access. /// public sealed class PerTenantWorkflowSource : IWorkflowSource { private readonly string _tenantDir; private readonly GitRunner _git; private readonly ILogger _logger; public PerTenantWorkflowSource(string tenantId, string tenantDir, ILogger logger, GitRunner git) { _tenantDir = tenantDir; _git = git; _logger = logger; _logger.LogInformation( "Per-tenant workflow source: tenant={TenantId} dir={Dir}", tenantId, _tenantDir); } public Task> ListAsync(CancellationToken ct) { ct.ThrowIfCancellationRequested(); EnsureInitialized(); return Task.FromResult(ListYamlFiles(_tenantDir)); } public async Task ReadAsync(string path, CancellationToken ct) { EnsureInitialized(); return await File.ReadAllTextAsync( Path.Combine(_tenantDir, path.Replace('/', Path.DirectorySeparatorChar)), ct); } public Task GetStateAsync(CancellationToken ct = default) { EnsureInitialized(); return _git.GetStateAsync(_tenantDir, ct); } /// /// Ensures the tenant directory exists. Workflows are strictly per-tenant: /// we never copy shared example templates into a tenant's directory, so a /// tenant sees only the definitions that belong to them. /// private void EnsureInitialized() { if (Directory.Exists(_tenantDir)) return; _logger.LogInformation("Initializing (empty) tenant workflow directory {Dir}", _tenantDir); try { Directory.CreateDirectory(_tenantDir); } catch (Exception ex) { throw new InvalidOperationException( $"Cannot create tenant workflow directory '{_tenantDir}'. " + "Check that WorkflowSource:CopiesRoot points to a writable location " + $"and that the process has filesystem permissions. {ex.Message}", ex); } } private static IReadOnlyList ListYamlFiles(string dir) { var files = new List(); if (!Directory.Exists(dir)) return files; foreach (var full in Directory.EnumerateFiles(dir, "*", SearchOption.AllDirectories)) { var ext = Path.GetExtension(full); if (ext is not (".yaml" or ".yml")) continue; // Skip .git var rel = Path.GetRelativePath(dir, full); if (rel.StartsWith(".git", StringComparison.OrdinalIgnoreCase)) continue; files.Add(new WorkflowFile(rel.Replace(Path.DirectorySeparatorChar, '/'))); } return files.OrderBy(f => f.Path, StringComparer.Ordinal).ToList(); } }