using System;
using System.Collections.Generic;
using System.Data.Common;
using System.Linq;
using System.Threading.Tasks;
using GB5Shared.DTO.Framework.Login;
using GB5Shared.DTO.WorkFlow;
using GB5Shared.Query.WorkFlow;
using GB5Shared.QueryExecutor;
using GB5Shared.Telemetry;
using GB5Shared.WorkFlow.WorkFlowRunTime;
using Microsoft.Extensions.Caching.Memory;
using Microsoft.Extensions.Logging;
using Newtonsoft.Json;
using Newtonsoft.Json.Linq;
using static GB5Shared.DTO.WorkFlow.ContractsDTO;
using static GB5Shared.GB5Constant.Constant;
namespace GB5Shared.WorkFlow.WorkFlowEngine
{
///
/// Dynamic, config-driven workflow engine.
///
/// Flow overview
/// ─────────────
/// 1.
/// → Queries MWORKFLOWCONFIG (EntityId / OUId / BizTransactionClassId /
/// BizTransactionId with wildcard -1 support).
/// → Evaluates the EVALCONDITION expression of each candidate row
/// using .
/// → Returns the first matching workflow configuration.
///
/// 2.
/// → Resolves the workflow definition from the config.
/// → Creates a TWORKFLOWINSTANCE row.
/// → Enters the first approval level (lowest APPROVALLEVEL number).
/// → For each step at that level, evaluates EVALCONDITION; creates
/// TWORKFLOWTASK rows only for steps where the condition is true.
/// → Resolves the assignee via the step's STRATEGYTYPE / RULEEXPRESSION.
///
/// 3.
/// → Marks the task complete.
/// → When ALL pending tasks at the current approval level are resolved:
/// - Approve → advance to next approval level (or complete).
/// - Reject → terminate the instance as Rejected.
/// - Return → go back to previous level (or to creator at level 0).
/// → Writes a TWORKFLOWHISTORY row for every action.
///
public class WorkFlowEngine : IWorkFlowEngine
{
private readonly IWorkFlowRunTime _runtime;
private readonly IMemoryCache _cache;
private readonly IUserDelegationResolver _delegation;
private readonly IQueryExecutor _qe;
private readonly ILogger _logger;
private static readonly TimeSpan CacheTtl = TimeSpan.FromMinutes(10);
public WorkFlowEngine(
IWorkFlowRunTime runtime,
IMemoryCache memoryCache,
IUserDelegationResolver delegationResolver,
IQueryExecutor queryExecutor,
ILogger logger)
{
_runtime = runtime;
_cache = memoryCache;
_delegation = delegationResolver;
_qe = queryExecutor;
_logger = logger;
}
// ══════════════════════════════════════════════════════════════════════
// AUTO-APPROVE — overdue task query
// ══════════════════════════════════════════════════════════════════════
public async Task> GetOverdueAutoApproveTasksAsync(
LoginDTO login,
DbTransaction tx)
{
return await _runtime.GetOverdueAutoApproveTasksAsync(login, tx);
}
// ══════════════════════════════════════════════════════════════════════
// 1. APPLICABILITY CHECK
// ══════════════════════════════════════════════════════════════════════
public async Task CheckWorkFlowApplicability(
WorkflowCheckContext context,
LoginDTO login,
DbTransaction tx)
{
try
{
ValidateCheckContext(context);
// ── WIP bypass: if this call is an approved WIP callback, skip all checks ──
if (context.Facts.TryGetValue(IWorkFlowEngine.WipApprovalFactKey, out var wipFlag)
&& wipFlag is bool b && b)
{
return new WorkflowPreCheckResult { IsEnabled = false };
}
// ── Step 1: Load MWORKFLOWCONFIG candidates ───────────────────
var configs = await _runtime.GetWorkflowConfigsAsync(
entityId: context.EntityId,
clientId: login.ClientId,
ouId: context.OUId,
bizTransactionClassId: context.BizTransactionClassId,
bizTransactionId: context.BizTransactionId,
loginDTO: login,
dbTransaction: tx);
if (configs == null || configs.Count == 0)
{
return new WorkflowPreCheckResult { IsEnabled = false };
}
// ── Step 2: Evaluate EVALCONDITION for each candidate ─────────
foreach (var config in configs)
{
bool conditionMet = EvaluateConfigCondition(config.EvalCondition, context.Facts);
if (!conditionMet)
continue;
// ── Step 3: Verify the referenced MWORKFLOW exists & active ─
// WorkflowId = -1 is the sentinel for "no workflow assigned" to this
// config row (same -1 convention used by ClientId/EntityId/OUId/etc.)
// — treat it as not applicable and keep looking, rather than erroring.
// WC.WORKFLOWID is read directly off MWORKFLOWCONFIG (no JOIN involved),
// so once -1 is ruled out every remaining value is a real, valid
// autonumber-generated WorkflowId — GB5 IDs are negative by convention
// (e.g. -1399999762), so there is no other invalid state to guard here.
if (config.WorkflowId == -1)
{
continue;
}
// ── Step 4: Determine mode — WIP or Regular ───────────────
int mode = config.IsFormBasedApproval
? WORKFLOW.WIP
: WORKFLOW.REGULAR;
return new WorkflowPreCheckResult
{
IsEnabled = true,
WorkFlowMode = mode,
WorkFlowId = config.WorkflowId,
WorkflowConfigId = config.WorkflowConfigId,
IsFormBasedApproval = config.IsFormBasedApproval,
CallbackEndpoint = config.CallbackEndpoint
};
}
// No config row's condition matched → no workflow
return new WorkflowPreCheckResult { IsEnabled = false };
}
catch (WorkflowEngineException)
{
throw;
}
catch (Exception ex)
{
throw new WorkflowEngineException(
"Workflow applicability check failed. " + ex.Message, ex);
}
}
// ══════════════════════════════════════════════════════════════════════
// 2. START WORKFLOW
// ══════════════════════════════════════════════════════════════════════
public async Task StartWorkflowAsync(
WorkflowStartRequest request,
LoginDTO login,
DbTransaction tx)
{
try
{
ValidateStartRequest(request);
// ── Step 1: Resolve applicable config ─────────────────────────
// Enrich facts with every scalar field from the posted DTO under
// the "DTO." prefix so that EVALCONDITION expressions like
// "DTO.TaskDetailType == 4" resolve correctly — both here and in
// EnterApprovalLevelAsync where step conditions are evaluated.
var enrichedFacts = new Dictionary(
request.Facts ?? new Dictionary(),
StringComparer.OrdinalIgnoreCase);
FlattenDataJsonFacts(request.DataJson, enrichedFacts);
var checkCtx = new WorkflowCheckContext
{
EntityId = request.EntityId,
OUId = request.OUId,
BizTransactionClassId = request.BizTransactionClassId,
BizTransactionId = request.BizTransactionId,
Facts = enrichedFacts
};
var preCheck = await CheckWorkFlowApplicability(checkCtx, login, tx);
if (!preCheck.IsEnabled)
return new WorkflowStartResult { IsStarted = false };
// ── Step 2: Load full workflow definition ─────────────────────
var def = await _runtime.GetWorkflowByIdAsync(preCheck.WorkFlowId, login, tx)
?? throw new WorkflowEngineException(
$"Workflow {preCheck.WorkFlowId} referenced in MWORKFLOWCONFIG " +
$"(ConfigId={preCheck.WorkflowConfigId}) was not found or is inactive.");
ValidateWorkflowDefinition(def);
int firstLevel = def.GetFirstApprovalLevel();
// ── Step 3: Resolve BizTransactionTypeId ──────────────────────
int bizTypeId = await _runtime.GetBizTransactionTypeIdAsync(
request.BizTransactionClassId, login, tx);
// ── WIP path: Insert TWORKFLOWWIP first, use WipId as ObjectId ─
int wipId = 0;
int objectId = request.ObjectId;
if (preCheck.IsFormBasedApproval)
{
string dataJsonStr = request.DataJson is string s
? s
: JsonConvert.SerializeObject(request.DataJson);
// Prefer explicit config endpoint; fall back to the URL of the current request
// so zero config is needed — the submitting endpoint IS the callback endpoint.
string? callbackEndpoint = !string.IsNullOrWhiteSpace(preCheck.CallbackEndpoint)
? preCheck.CallbackEndpoint
: login.RequestUrl;
// Serialize the submitter's full login so the dispatcher can replay
// the entity save under the original session context (WorkOUId, BranchId,
// UserId, etc.) rather than the approver's context.
// WipApprovalId and RequestUrl are [JsonIgnore] — never serialized.
string? submitterLoginJson = JsonConvert.SerializeObject(login);
wipId = await _runtime.InsertWipAsync(
entityId: request.EntityId,
tenantId: login.ClientId,
dataJson: dataJsonStr,
loginJson: submitterLoginJson,
apiEndpoint: callbackEndpoint,
createdById: login.UserId,
loginDTO: login,
dbTransaction: tx);
objectId = wipId; // WIP ID becomes the instance's ObjectId
}
// ── Step 3b: Remove any stale pending instance for this entity+object ──
// Prevents duplicate pending rows (e.g. from a double-submit or a
// previous failed run) that would block the approval inbox query.
await _runtime.DeleteDuplicateInstanceAsync(
request.EntityId, objectId, login, tx);
// ── Step 4: Create workflow instance ──────────────────────────
var instance = new WorkflowInstance
{
EntityId = request.EntityId,
WorkflowId = def.WorkflowId,
ObjectId = objectId,
BizTransactionTypeId = bizTypeId,
DataJson = request.DataJson,
// Persist the full enriched facts so every approval level can resolve
// EnrichQualifier (StrategyType=2) assignments and condition expressions
// without the caller re-supplying them.
//
// enrichedFacts = request.Facts (DB_ENRICH bag, e.g. ReportingToEmployeeId)
// + all scalar DTO fields as "DTO." keys
// (added by FlattenDataJsonFacts above).
//
// Using request.Facts alone caused FACTSJSON = NULL whenever no DB_ENRICH
// qualifier was configured for the entity — the bag was empty and
// SerializeFactsJson returned null. enrichedFacts is always populated
// as long as the DTO has scalar fields, so FACTSJSON is never null.
//
// BuildFactsFromDataJson (called at action time) layers FactsJson (Layer 1)
// under DataJson (Layer 2), so DTO field duplication is harmless —
// DataJson simply overwrites the same values at higher priority.
FactsJson = SerializeFactsJson(enrichedFacts),
WorkflowStatus = (int)WorkflowInstanceStatus.Pending,
CurrentApprovalLevel = firstLevel,
CurrentStepId = -1
};
instance = await _runtime.InsertInstanceAsync(instance, login, tx);
// ── WIP: bind instance ID back to the WIP row ─────────────────
// wipId defaults to the local 0 above and is only reassigned via
// InsertWipAsync's real (negative autonumber) row id — so "was a WIP
// actually created" must be tested as != 0, never > 0.
if (preCheck.IsFormBasedApproval && wipId != 0)
{
await _runtime.UpdateWipInstanceIdAsync(
wipId, instance.WorkflowInstanceId, login.UserId, login, tx);
}
// ── Step 5: Record initial history ────────────────────────────
await _runtime.InsertHistoryAsync(new WorkflowHistory
{
InstanceId = instance.WorkflowInstanceId,
WorkflowId = def.WorkflowId,
EntityId = instance.EntityId,
ObjectId = objectId,
BiztransactionTypeId = bizTypeId,
DataJson = request.DataJson,
StepKey = "SUBMIT",
ApprovalLevel = 0,
Action = (int)WorkflowAction.Submit,
ActionByUserId = login.UserId,
Comment = request.Comment,
ActionOn = DateTime.UtcNow
}, login, tx);
// ── Step 6: Enter first approval level ────────────────────────
// earlyDispatch is non-null only for the edge case where the first level
// has no applicable steps and the WIP finalises immediately at submission.
var pendingNotifications = new List();
var (earlyDispatch, _) = await EnterApprovalLevelAsync(
instance, def, firstLevel, enrichedFacts, login, tx, pendingNotifications);
return new WorkflowStartResult
{
IsStarted = true,
IsWip = preCheck.IsFormBasedApproval,
WipId = wipId,
WorkflowInstanceId = instance.WorkflowInstanceId,
PendingNotifications = pendingNotifications,
PendingDispatch = earlyDispatch
};
}
catch (WorkflowEngineException)
{
throw;
}
catch (Exception ex)
{
throw new WorkflowEngineException(
$"StartWorkflowAsync failed for Entity={request.EntityId}, " +
$"Object={request.ObjectId}. {ex.Message}", ex);
}
}
// ══════════════════════════════════════════════════════════════════════
// 2b. CANCEL ON DELETE / CANCEL ON RESUBMIT
// ══════════════════════════════════════════════════════════════════════
public async Task CancelPendingInstanceAsync(
int entityId,
int objectId,
string reason,
LoginDTO login,
DbTransaction tx)
{
try
{
var cancelled = await _runtime.CancelPendingInstancesAsync(entityId, objectId, login, tx);
foreach (var instance in cancelled)
{
await _runtime.InsertHistoryAsync(new WorkflowHistory
{
InstanceId = instance.WorkflowInstanceId,
WorkflowId = instance.WorkflowId,
EntityId = entityId,
ObjectId = objectId,
BiztransactionTypeId = instance.BizTransactionTypeId,
DataJson = instance.DataJson,
StepKey = "CANCEL",
ApprovalLevel = 0,
Action = (int)WorkflowAction.Cancel,
ActionByUserId = login.UserId,
Comment = reason,
ActionOn = DateTime.UtcNow
}, login, tx);
}
return cancelled.Count;
}
catch (Exception ex)
{
throw new WorkflowEngineException(
$"CancelPendingInstanceAsync failed for Entity={entityId}, Object={objectId}. {ex.Message}", ex);
}
}
// ══════════════════════════════════════════════════════════════════════
// 3. ACTION HANDLING
// ══════════════════════════════════════════════════════════════════════
public Task UpdateWipObjectIdAsync(int wipId, int objectId, LoginDTO login, DbTransaction tx)
=> _runtime.UpdateWipStatusAsync(wipId, 2, 1, objectId, login.UserId, login, tx);
public async Task HandleActionsAsync(
IEnumerable actions,
LoginDTO login,
DbTransaction tx)
{
var result = new WorkflowActionResult();
foreach (var action in actions)
{
var (dispatch, outcome, item) = await HandleSingleActionAsync(action, login, tx);
if (dispatch != null) result.PendingDispatches.Add(dispatch);
if (outcome != null) result.Outcomes.Add(outcome);
if (item != null) result.ProcessedItems.Add(item);
}
return result;
}
private async Task<(WipDispatchInfo? Dispatch, WorkflowOutcomeInfo? Outcome, WorkflowActionItemInfo? Item)> HandleSingleActionAsync(
WorkflowActionContext ctx,
LoginDTO login,
DbTransaction tx)
{
// ── Load & validate task ──────────────────────────────────────────
var task = await _runtime.GetTaskAsync(ctx.TaskId, login, tx)
?? throw new WorkflowEngineException(
$"Workflow task {ctx.TaskId} was not found.");
if ((WorkflowTaskStatus)task.WorkflowTaskStatus != WorkflowTaskStatus.Pending)
{
throw new WorkflowEngineException(
$"Task {ctx.TaskId} is no longer pending " +
$"(current status: {(WorkflowTaskStatus)task.WorkflowTaskStatus}). " +
"Action cannot be performed on a completed or rejected task.");
}
var instance = await _runtime.GetInstanceAsync(task.WorkflowInstanceId, login, tx)
?? throw new WorkflowEngineException(
$"Workflow instance {task.WorkflowInstanceId} not found.");
// DATAJSON may be null for instances created before the column was added to the schema,
// or for rows inserted by an older version of the engine that omitted it.
// Fall back to the WIP row's DataJson — it always carries the original DTO snapshot.
bool instanceDataJsonMissing =
instance.DataJson is null ||
(instance.DataJson is string ds && string.IsNullOrEmpty(ds));
if (instanceDataJsonMissing)
{
var wipForFacts = await _runtime.GetWipByInstanceIdAsync(
instance.WorkflowInstanceId, login, tx);
if (wipForFacts?.DataJson != null)
instance.DataJson = wipForFacts.DataJson;
}
var def = await _runtime.GetWorkflowByIdAsync(instance.WorkflowId, login, tx)
?? throw new WorkflowEngineException(
$"Workflow definition {instance.WorkflowId} not found or inactive.");
ValidateWorkflowDefinition(def);
// Captured once, up front, so it's available on every return path below —
// including the "still pending at level" case, which never yields an Outcome.
var item = new WorkflowActionItemInfo
{
TaskId = ctx.TaskId,
WorkflowInstanceId = instance.WorkflowInstanceId,
EntityId = instance.EntityId,
ObjectId = instance.ObjectId,
Action = ctx.Action,
DataJson = instance.DataJson
};
// ── Find the step this task belongs to ────────────────────────────
var step = def.Steps.FirstOrDefault(s => s.WorkflowDetailId == task.StepId)
?? throw new WorkflowEngineException(
$"Step {task.StepId} does not exist in workflow {def.WorkflowId}.");
int currentLevel = instance.CurrentApprovalLevel;
// ── Complete the task ─────────────────────────────────────────────
task.WorkflowTaskStatus = (int)WorkflowTaskStatus.Completed;
task.ActionTaken = ctx.Action;
task.Comment = ctx.Comment;
task.CompletedOn = DateTime.UtcNow;
await _runtime.UpdateTaskAsync(task, login, tx);
// ── Write history ─────────────────────────────────────────────────
await _runtime.InsertHistoryAsync(new WorkflowHistory
{
InstanceId = instance.WorkflowInstanceId,
WorkflowId = def.WorkflowId,
EntityId = instance.EntityId,
ObjectId = instance.ObjectId,
StepKey = step.StepKey,
ApprovalLevel = currentLevel,
Action = ctx.Action,
ActionByUserId = ctx.UserId,
Comment = ctx.Comment,
ActionOn = DateTime.UtcNow
}, login, tx);
// ── Handle REJECT — terminate immediately ─────────────────────────
if (ctx.Action == (int)WorkflowAction.Reject)
{
instance.WorkflowStatus = (int)WorkflowInstanceStatus.Rejected;
await _runtime.UpdateInstanceAsync(instance, login, tx);
// Update WIP status if applicable; otherwise update entity STATUS = 3
var wipOnReject = await _runtime.GetWipByInstanceIdAsync(instance.WorkflowInstanceId, login, tx);
if (wipOnReject != null)
await _runtime.UpdateWipStatusAsync(wipOnReject.WipId, 3, 2, 0, login.UserId, login, tx);
else
await SetEntityStatusAsync(instance, 3, login, tx);
return (null, new WorkflowOutcomeInfo
{
WorkflowInstanceId = instance.WorkflowInstanceId,
WorkflowId = instance.WorkflowId,
EntityId = instance.EntityId,
ObjectId = instance.ObjectId,
Action = ctx.Action,
ActionByUserId = ctx.UserId,
Comment = ctx.Comment,
DataJson = instance.DataJson,
FactsJson = instance.FactsJson
}, item);
}
// ── Rebuild facts from stored DataJson + FactsJson so level conditions and
// EnrichQualifier (StrategyType=2) assignments work at every level without
// the caller re-supplying them. ctx.Facts are merged last as overrides.
var instanceFacts = BuildFactsFromDataJson(instance.DataJson, instance.FactsJson, ctx.Facts);
// ── Handle RETURN — go back to previous level ─────────────────────
if (ctx.Action == (int)WorkflowAction.Return)
{
var levels = def.GetApprovalLevels();
int currentIdx = levels.IndexOf(currentLevel);
int prevLevel = currentIdx > 0 ? levels[currentIdx - 1] : 0;
instance.WorkflowStatus = (int)WorkflowInstanceStatus.Returned;
instance.CurrentApprovalLevel = prevLevel;
await _runtime.UpdateInstanceAsync(instance, login, tx);
// Update WIP status if applicable
var wipOnReturn = await _runtime.GetWipByInstanceIdAsync(instance.WorkflowInstanceId, login, tx);
if (wipOnReturn != null)
await _runtime.UpdateWipStatusAsync(wipOnReturn.WipId, 4, 3, 0, login.UserId, login, tx);
WipDispatchInfo? returnDispatch = null;
if (prevLevel > 0)
{
(returnDispatch, _) = await EnterApprovalLevelAsync(
instance, def, prevLevel, instanceFacts, login, tx);
}
return (returnDispatch, new WorkflowOutcomeInfo
{
WorkflowInstanceId = instance.WorkflowInstanceId,
WorkflowId = instance.WorkflowId,
EntityId = instance.EntityId,
ObjectId = instance.ObjectId,
Action = ctx.Action,
ActionByUserId = ctx.UserId,
Comment = ctx.Comment,
DataJson = instance.DataJson,
FactsJson = instance.FactsJson
}, item);
}
// ── APPROVE — check whether all tasks at current level are done ───
var pendingAtLevel = await _runtime.GetPendingTasksForLevelAsync(
instance.WorkflowInstanceId, currentLevel, login, tx);
if (pendingAtLevel.Any())
{
// Other approvers at this level are still pending — wait for them
return (null, null, item);
}
// ── All tasks at current level approved → advance ─────────────────
int? nextLevel = def.GetNextApprovalLevel(currentLevel);
if (nextLevel.HasValue)
{
int levelToEnter = nextLevel.Value;
// Same-approver auto-skip: if a subsequent level's resolved approver(s) are
// identical to the level just completed, that level is redundant (the same
// person would just be approving their own approval again) — skip it, record
// an auto-approved history row, and check the level after it, cascading
// forward until a level with a genuinely different approver is found or the
// workflow runs out of levels.
var anchor = await ResolveLevelApproverIdentityAsync(def, currentLevel, instanceFacts, login, tx);
if (anchor.IsComparable)
{
while (true)
{
var candidate = await ResolveLevelApproverIdentityAsync(def, levelToEnter, instanceFacts, login, tx);
bool sameApprover = candidate.IsComparable && candidate.UserIds.SetEquals(anchor.UserIds);
if (!sameApprover)
break;
var skippedStepKeys = string.Join(",", def.GetStepsAtLevel(levelToEnter).Select(s => s.StepKey));
await _runtime.InsertHistoryAsync(new WorkflowHistory
{
InstanceId = instance.WorkflowInstanceId,
WorkflowId = def.WorkflowId,
EntityId = instance.EntityId,
ObjectId = instance.ObjectId,
StepKey = skippedStepKeys,
ApprovalLevel = levelToEnter,
Action = (int)WorkflowAction.Approve,
ActionByUserId = ctx.UserId,
Comment = "Auto-approved: same approver as previous level.",
ActionOn = DateTime.UtcNow
}, login, tx);
var afterSkip = def.GetNextApprovalLevel(levelToEnter);
if (!afterSkip.HasValue)
{
var (skipDispatch, skipOutcome) = await FinaliseCompletionAsync(instance, login, tx);
return (skipDispatch, skipOutcome, item);
}
levelToEnter = afterSkip.Value;
}
}
instance.CurrentApprovalLevel = levelToEnter;
instance.WorkflowStatus = (int)WorkflowInstanceStatus.Pending;
await _runtime.UpdateInstanceAsync(instance, login, tx);
// Capture dispatch and outcome: if the next level has no applicable steps
// it auto-finalises immediately — both must reach the BLL so the outcome
// event (APPROVEWORKFLOW) is published even in the auto-skip case.
var (levelDispatch, levelOutcome) = await EnterApprovalLevelAsync(
instance, def, levelToEnter, instanceFacts, login, tx);
return (levelDispatch, levelOutcome, item);
}
else
{
// Final level approved → workflow complete
var (finalDispatch, finalOutcome) = await FinaliseCompletionAsync(instance, login, tx);
return (finalDispatch, finalOutcome, item);
}
}
///
/// Resolves the set of concrete UserIds who would approve the applicable steps at
/// , used to detect "this level's approver is the
/// same person as the previous level's approver" for the auto-skip feature.
/// is false whenever any applicable step resolves to
/// a Role/Group pool assignee (StrategyType=UserGroup, or any assignee lacking a
/// concrete UserId) — "same approver" has no well-defined meaning for a pool of
/// people, so such levels are never auto-skipped.
///
/// Note: a concrete UserId here can legitimately be negative — this system generates
/// UserIds via autonumber (e.g. RULEEXPRESSION = '-1399999762' in
/// 's StrategyType=User case) — so only the exact
/// sentinel value -1 (meaning "no approver resolved") disqualifies comparability, not
/// "any non-positive UserId".
///
private async Task<(HashSet UserIds, bool IsComparable)> ResolveLevelApproverIdentityAsync(
WorkflowDefinition def,
int approvalLevel,
IReadOnlyDictionary facts,
LoginDTO login,
DbTransaction tx)
{
var userIds = new HashSet();
foreach (var step in def.GetStepsAtLevel(approvalLevel))
{
if (!EvaluateStepCondition(step.EvalCondition, facts, step))
continue;
var assignees = await ResolveAssigneesAsync(step, def.EntityId, facts, login, tx);
foreach (var assignee in assignees)
{
if (assignee.UserId == -1 || assignee.RoleId != -1 || assignee.GroupId != -1)
return (userIds, false);
userIds.Add(assignee.UserId);
}
}
return (userIds, userIds.Count > 0);
}
// ══════════════════════════════b════════════════════════════════════════
// LEVEL ENTRY — evaluate conditions, create tasks
// ══════════════════════════════════════════════════════════════════════
///
/// Enters an approval level:
/// 1. Evaluates EVALCONDITION for every step at .
/// 2. Creates TWORKFLOWTASK rows for steps whose condition evaluates to true.
/// 3. If steps exist at this level but NO condition matches → workflow is
/// finalised as Completed and the main object status is updated.
/// 4. If no steps are defined at this level at all → advance to next level.
///
///
/// When non-null, one is appended
/// per task inserted. Callers that do not need notification data (e.g. action
/// Return/Advance paths) may pass null to skip collection.
///
///
/// A when this level entry immediately finalises the
/// workflow (no steps configured, or no step conditions matched) and the instance
/// is a WIP — so the caller can propagate the dispatch after committing the transaction.
/// Returns null for regular (non-WIP) instances, or when approval tasks were
/// created and the workflow is still pending.
///
private async Task<(WipDispatchInfo? Dispatch, WorkflowOutcomeInfo? Outcome)> EnterApprovalLevelAsync(
WorkflowInstance instance,
WorkflowDefinition def,
int approvalLevel,
IReadOnlyDictionary facts,
LoginDTO login,
DbTransaction tx,
List? pendingNotifications = null)
{
var stepsAtLevel = def.GetStepsAtLevel(approvalLevel);
if (!stepsAtLevel.Any())
{
// No steps configured for this level — advance to next or finalise
int? next = def.GetNextApprovalLevel(approvalLevel);
if (next.HasValue)
return await EnterApprovalLevelAsync(instance, def, next.Value, facts, login, tx, pendingNotifications);
return await FinaliseCompletionAsync(instance, login, tx);
}
// ── Evaluate EVALCONDITION for every step at this level ───────────
var applicableSteps = new List();
foreach (var step in stepsAtLevel)
{
if (EvaluateStepCondition(step.EvalCondition, facts, step))
applicableSteps.Add(step);
}
if (!applicableSteps.Any())
{
// Steps exist but none of their conditions matched the current document.
// No approval is required at this level → finalise workflow immediately.
return await FinaliseCompletionAsync(instance, login, tx);
}
// ── Create one task per applicable step ───────────────────────────
DateTime now = DateTime.UtcNow;
bool anyTaskCreated = false;
foreach (var step in applicableSteps)
{
var assignees = await ResolveAssigneesAsync(step, instance.EntityId, facts, login, tx);
// "No approver" sentinel: UserId=-1 with no Role/Group means the
// assignment rule (e.g. an EnrichQualifier walking a reporting hierarchy,
// or a static rule intentionally configured with RULEEXPRESSION=-1) found
// nobody to approve at this step — there's nothing above this point in the
// chain. Treat it as "this step needs no approval": never create a task for
// it (a task assigned to UserId=-1 could never be actioned by anyone and
// would hang the workflow forever) — record it as auto-approved instead.
// This applies uniformly whether the workflow has one level, two levels, or
// many: whichever level's approver resolves to -1 is the final level that
// actually required a person, and it auto-completes right there.
var realAssignees = assignees
.Where(a => !(a.UserId == -1 && a.RoleId == -1 && a.GroupId == -1))
.ToList();
if (realAssignees.Count < assignees.Count)
{
await _runtime.InsertHistoryAsync(new WorkflowHistory
{
InstanceId = instance.WorkflowInstanceId,
WorkflowId = def.WorkflowId,
EntityId = instance.EntityId,
ObjectId = instance.ObjectId,
StepKey = step.StepKey,
ApprovalLevel = approvalLevel,
Action = (int)WorkflowAction.Approve,
ActionByUserId = login.UserId,
Comment = "Auto-approved: no approver resolved (Id=-1).",
ActionOn = now
}, login, tx);
}
foreach (var assignee in realAssignees)
{
// A real GB5 UserId is a negative autonumber value (e.g. -1399999762),
// so "does this assignee have a concrete UserId to delegate-resolve"
// must be assignee.UserId != -1 (the unset sentinel), never > 0 — a
// UserGroup assignee correctly leaves UserId at -1 (unset) and skips
// delegation resolution, exactly as before.
int effectiveUserId = assignee.UserId != -1
? await ResolveEffectiveUserAsync(assignee.UserId, now, login, tx)
: assignee.UserId;
var insertedTask = await _runtime.InsertTaskAsync(new WorkflowTask
{
WorkflowInstanceId = instance.WorkflowInstanceId,
StepId = step.WorkflowDetailId,
WorkflowTaskStatus = (int)WorkflowTaskStatus.Pending,
AssignedToUserId = effectiveUserId,
AssignedRoleId = assignee.RoleId,
AssignedUserGroupId = assignee.GroupId,
DueOn = step.SlaHours > 0
? now.AddMinutes(step.SlaHours)
: (DateTime?)null,
CreatedOn = now
}, login, tx);
instance.CurrentStepId = step.WorkflowDetailId;
anyTaskCreated = true;
pendingNotifications?.Add(new WorkflowTaskNotificationRequest
{
TaskId = insertedTask.WorkflowTaskId,
AssigneeUserId = effectiveUserId,
WorkflowInstanceId = instance.WorkflowInstanceId,
WorkflowId = def.WorkflowId,
DataJson = instance.DataJson,
FactsJson = instance.FactsJson
});
}
}
if (!anyTaskCreated)
{
// Every applicable step at this level resolved to "no approver" (-1) —
// this level is complete with nothing to wait on. Cascade forward exactly
// like an unconfigured level: advance to the next level, or finalise if
// this was the last one.
int? nextLevelAfterAutoApprove = def.GetNextApprovalLevel(approvalLevel);
if (nextLevelAfterAutoApprove.HasValue)
return await EnterApprovalLevelAsync(
instance, def, nextLevelAfterAutoApprove.Value, facts, login, tx, pendingNotifications);
return await FinaliseCompletionAsync(instance, login, tx);
}
await _runtime.UpdateInstanceAsync(instance, login, tx);
return (null, null); // approval tasks created — workflow is pending, no immediate dispatch
}
// ══════════════════════════════════════════════════════════════════════
// ASSIGNEE RESOLUTION
// ══════════════════════════════════════════════════════════════════════
private async Task> ResolveAssigneesAsync(
WorkflowStep step,
int entityId,
IReadOnlyDictionary facts,
LoginDTO login,
DbTransaction tx)
{
var results = new List();
// Guard: MWORKFLOWDETAIL.ASSIGNMENTID is a FK to MWORKFLOWASSIGNMENT.
// Dapper maps a LEFT JOIN miss (SQL NULL) to this non-nullable int's default,
// 0 — the exact and only "no assignment" sentinel here. A real matched row's
// WorkflowAssignmentId is a negative GB5 autonumber value (e.g. -1399999762),
// so this must be == 0, never <= 0 (that would misfire on every real match).
if (step.WorkflowAssignmentId == 0)
throw new WorkflowEngineException(
$"Step '{step.DisplayName}' (Id={step.WorkflowDetailId}) has " +
$"ASSIGNMENTID={step.AssignmentId} but no matching row exists in " +
"MWORKFLOWASSIGNMENT. " +
"Insert a MWORKFLOWASSIGNMENT row with that WORKFLOWASSIGNMENTID and " +
"set STRATEGYTYPE + RULEEXPRESSION correctly.");
switch (step.StrategyType)
{
// ── 0: Static User ────────────────────────────────────────────
// MWORKFLOWASSIGNMENT.RULEEXPRESSION holds the UserId as a string.
// e.g. RULEEXPRESSION = '-1399999762'
case 0:
{
if (!int.TryParse(step.RuleExpression, out int userId))
throw new WorkflowEngineException(
$"Step '{step.DisplayName}' (Id={step.WorkflowDetailId}) has " +
$"StrategyType=User but MWORKFLOWASSIGNMENT (Id={step.WorkflowAssignmentId}) " +
$"RULEEXPRESSION '{step.RuleExpression}' is not a valid UserId integer. " +
"Set RULEEXPRESSION to the approver's UserId.");
results.Add(new AssigneeResult { UserId = userId });
break;
}
// ── 1: UserGroup ──────────────────────────────────────────────
// MWORKFLOWASSIGNMENT.RULEEXPRESSION holds the UserGroupId as a string.
case 1:
{
if (!int.TryParse(step.RuleExpression, out int groupId))
throw new WorkflowEngineException(
$"Step '{step.DisplayName}' (Id={step.WorkflowDetailId}) has " +
$"StrategyType=UserGroup but MWORKFLOWASSIGNMENT (Id={step.WorkflowAssignmentId}) " +
$"RULEEXPRESSION '{step.RuleExpression}' is not a valid UserGroupId integer. " +
"Set RULEEXPRESSION to the UserGroupId.");
results.Add(new AssigneeResult { GroupId = groupId });
break;
}
// ── 2: EnrichQualifier ───────────────────────────────────────
// MWORKFLOWASSIGNMENT.RULEEXPRESSION holds the name of the field
// (e.g. "ReportingToEmployeeId") that a DB_ENRICH qualifier wrote
// into ctx._bag during EnrichForWorkflowAsync. That enriched value
// is passed here via WorkflowStartRequest.Facts (merged from the bag
// in EventHandler.ExecuteSaveAsync before calling StartWorkflowAsync).
// No additional DB query is required — the value is already in facts.
case 2:
{
var fieldName = step.RuleExpression?.Trim();
if (string.IsNullOrEmpty(fieldName))
throw new WorkflowEngineException(
$"Step '{step.DisplayName}' (Id={step.WorkflowDetailId}) has " +
"StrategyType=EnrichQualifier but RULEEXPRESSION is empty. " +
"Set RULEEXPRESSION to the field name written by the DB_ENRICH qualifier " +
"(e.g. 'ReportingToEmployeeId').");
if (!facts.TryGetValue(fieldName, out var rawValue) || rawValue is null)
throw await BuildMissingEnrichFactExceptionAsync(step, fieldName, entityId, login, tx);
if (!int.TryParse(rawValue.ToString(), out int enrichedUserId))
throw new WorkflowEngineException(
$"Step '{step.DisplayName}' (Id={step.WorkflowDetailId}) has " +
$"StrategyType=EnrichQualifier: field '{fieldName}' resolved to " +
$"'{rawValue}' which is not a valid positive UserId integer.");
results.Add(new AssigneeResult { UserId = enrichedUserId });
break;
}
// ── 3: Dynamic SQL ────────────────────────────────────────────
case 3:
{
if (string.IsNullOrWhiteSpace(step.RuleExpression))
throw new WorkflowEngineException(
$"Step '{step.DisplayName}' (Id={step.WorkflowDetailId}) has " +
"StrategyType=DynamicSQL but RuleExpression is empty. " +
"Please configure a valid SQL query in MWORKFLOWASSIGNMENT.");
// Convert facts to a Dapper parameter object
// Facts keys are passed as @FieldName parameters
var sqlParams = BuildSqlParameters(facts);
IEnumerable rows;
try
{
rows = await _qe.QueryAsync(
login, step.RuleExpression, sqlParams, tx);
}
catch (Exception ex)
{
throw new WorkflowEngineException(
$"Dynamic SQL for step '{step.DisplayName}' " +
$"(Id={step.WorkflowDetailId}) failed to execute. " +
$"SQL: [{step.RuleExpression}]. Error: {ex.Message}", ex);
}
var rowList = rows?.ToList();
if (rowList == null || !rowList.Any())
throw new WorkflowEngineException(
$"Dynamic SQL for step '{step.DisplayName}' " +
$"(Id={step.WorkflowDetailId}) returned no rows. " +
"Cannot determine assignee. " +
"Verify the SQL query and the facts passed to the workflow.");
// Dynamic SQL always resolves to exactly ONE assignee — the first row.
// If the SQL returns multiple rows (e.g. a broad WHERE clause),
// only the first is used. For multi-approver scenarios use
// parallel steps at the same approval level.
var firstRow = (IDictionary)rowList[0];
var firstEntry = firstRow.FirstOrDefault();
if (firstEntry.Value == null)
throw new WorkflowEngineException(
$"Dynamic SQL for step '{step.DisplayName}' " +
$"(Id={step.WorkflowDetailId}) returned a row where " +
$"the first column ('{firstEntry.Key}') is NULL. " +
"The query must return a valid UserId as its first column.");
if (!int.TryParse(firstEntry.Value.ToString(), out int uid))
throw new WorkflowEngineException(
$"Dynamic SQL for step '{step.DisplayName}' " +
$"(Id={step.WorkflowDetailId}) returned a non-integer value " +
$"'{firstEntry.Value}' in column '{firstEntry.Key}'. " +
"The first column must be a UserId (integer).");
results.Add(new AssigneeResult { UserId = uid });
break;
}
// ── 4: External service (reserved) ────────────────────────────
case 4:
throw new WorkflowEngineException(
$"Step '{step.DisplayName}' (Id={step.WorkflowDetailId}) uses " +
"StrategyType=Service which is not yet implemented.");
default:
throw new WorkflowEngineException(
$"Step '{step.DisplayName}' (Id={step.WorkflowDetailId}) has " +
$"unknown StrategyType={step.StrategyType}. " +
"Valid values: 0=User, 1=UserGroup, 2=EnrichQualifier, 3=DynamicSQL, 4=Service.");
}
return results;
}
// ══════════════════════════════════════════════════════════════════════
// ENRICH-QUALIFIER DIAGNOSTICS
//
// Builds a self-diagnosing WorkflowEngineException for the StrategyType=2
// (EnrichQualifier) "field not found in facts" case. The generic message alone
// sent an operator hunting through the Qualifier engine's five tables
// (MQUALIFIERENTITY / MQUALIFIERDEFINITION / MQUALIFIERVERSION / MQUALIFIERMETHOD /
// MGBCONFIGUREDQUERY) to find the actual gap. This probes the one query that
// distinguishes the three real causes and puts the answer directly in the
// exception message, which lands verbatim in the endpoint span's "error" /
// "gb5.app.general_errors" tag — visible in Zipkin without opening any code.
// Also emits a GB5Trace step so the check is visible on the span timeline even
// when this exception is caught/wrapped upstream and its message is lost.
// ══════════════════════════════════════════════════════════════════════
private const string DiagnoseEnrichBindingSql = @"
SELECT
e.ENTITYCODE AS EntityCode,
qe.QUALIFIERID AS QualifierId,
qe.ISACTIVE AS IsActiveRaw,
qd.QUALIFIERCODE AS QualifierCode,
(SELECT COUNT(*) FROM MQUALIFIERVERSION qv
WHERE qv.QUALIFIERID = qe.QUALIFIERID
AND qv.TENANTID = qe.TENANTID
AND qv.VERSIONSTATUS = 1) AS ActiveVersionCount
FROM MENTITY e
LEFT JOIN MQUALIFIERENTITY qe ON qe.ENTITYCODE = e.ENTITYCODE
AND qe.TENANTID = @TenantId
AND qe.STAGE = 3 -- PreWorkflow
AND qe.SCOPE = 1 -- Workflow
LEFT JOIN MQUALIFIERDEFINITION qd ON qd.QUALIFIERID = qe.QUALIFIERID
WHERE e.ENTITYID = @EntityId;";
private async Task BuildMissingEnrichFactExceptionAsync(
WorkflowStep step,
string fieldName,
int entityId,
LoginDTO login,
DbTransaction tx)
{
string baseMessage =
$"Step '{step.DisplayName}' (Id={step.WorkflowDetailId}) has " +
$"StrategyType=EnrichQualifier and expects field '{fieldName}' in facts, " +
"but it was not found or is null.";
string diagnosis;
string entityCode = "(unknown)";
int bindingCount = 0;
try
{
// Dictionary-cast on the dynamic Dapper rows, matching this file's existing
// DynamicSQL-case convention (see the StrategyType=3 block above) rather than
// dynamic property access.
var rawRows = await _qe.QueryAsync(
login, DiagnoseEnrichBindingSql,
new { EntityId = entityId, TenantId = login.ClientId }, tx);
var rows = rawRows.Select(r => (IDictionary)r).ToList();
if (rows.Count > 0 && rows[0]["EntityCode"] is string ec)
entityCode = ec;
var bindings = rows.Where(r => r["QualifierId"] is not null).ToList();
bindingCount = bindings.Count;
if (bindings.Count == 0)
{
diagnosis =
$"DIAGNOSIS: No MQUALIFIERENTITY binding exists for EntityCode='{entityCode}', " +
"Stage=3(PreWorkflow), Scope=1(Workflow), TenantId=" + login.ClientId + ". " +
$"FIX: create a Qualifier (QUALIFIERTYPE=2/Enrich) whose query/method outputs a " +
$"fact column named exactly '{fieldName}', bind it in MQUALIFIERENTITY " +
$"(EntityCode='{entityCode}', Stage=3, Scope=1, ISACTIVE=0), and give it an Active " +
"MQUALIFIERVERSION (a bound qualifier with no active compiled version never runs).";
}
else
{
var parts = new List();
foreach (var b in bindings)
{
bool isActive = Convert.ToByte(b["IsActiveRaw"]) == 0; // ISACTIVE is inverted: 0=Active
int activeVersions = Convert.ToInt32(b["ActiveVersionCount"]);
string qCode = (string)b["QualifierCode"];
int qId = Convert.ToInt32(b["QualifierId"]);
if (!isActive)
parts.Add(
$"Qualifier '{qCode}' (Id={qId}) is bound to this entity/stage/scope but " +
"INACTIVE (MQUALIFIERENTITY.ISACTIVE=1). FIX: set ISACTIVE=0 on that row.");
else if (activeVersions == 0)
parts.Add(
$"Qualifier '{qCode}' (Id={qId}) is bound and active, but has NO Active " +
"MQUALIFIERVERSION — it is configured but never executes. FIX: compile and " +
$"activate a version for QualifierId={qId} (MQUALIFIERVERSION.VERSIONSTATUS=1 " +
"with a valid SNAPSHOTJSON).");
else
parts.Add(
$"Qualifier '{qCode}' (Id={qId}) is bound, active, and has an active compiled " +
$"version, but did not produce a fact named '{fieldName}' for this document. " +
"FIX: check its MQUALIFIERMETHOD -> MGBCONFIGUREDQUERY — the query must return " +
$"a column aliased exactly '{fieldName}' (case-sensitive) and must return a row " +
"for this document's parameters (a query with no matching row writes no facts).");
}
diagnosis = "DIAGNOSIS: " + string.Join(" | ", parts);
}
}
catch (Exception diagEx)
{
// Diagnostics must never mask or replace the real error — fall back to
// the plain message if the probe itself fails for any reason.
diagnosis = $"DIAGNOSIS: probe failed ({diagEx.Message}) — check MQUALIFIERENTITY/" +
"MQUALIFIERVERSION manually for EntityId=" + entityId + ", Stage=3, Scope=1.";
}
GB5Trace.Step("enrich-qualifier-diagnosis", new
{
stepId = step.WorkflowDetailId,
stepName = step.DisplayName,
fieldName,
entityCode,
tenantId = login.ClientId,
bindingCount,
diagnosis
});
return new WorkflowEngineException($"{baseMessage} {diagnosis}");
}
// ══════════════════════════════════════════════════════════════════════
// DELEGATION
// ══════════════════════════════════════════════════════════════════════
public async Task ResolveEffectiveUserAsync(
long userId,
DateTime at,
LoginDTO login,
DbTransaction tx)
{
try
{
var delegations = await _delegation.GetAsync(userId, login, tx);
var active = delegations.FirstOrDefault(
d => d.StartsOn <= at && d.EndsOn >= at);
return (int)(active?.DelegateUserId ?? userId);
}
catch (Exception) { throw; }
}
// ══════════════════════════════════════════════════════════════════════
// CONDITION EVALUATION
// ══════════════════════════════════════════════════════════════════════
///
/// Evaluates the EVALCONDITION of a MWORKFLOWCONFIG row against the facts.
/// Returns true when expression is null/empty (no restriction).
///
private bool EvaluateConfigCondition(
string? expression,
IReadOnlyDictionary facts)
{
if (string.IsNullOrWhiteSpace(expression))
return true;
try
{
return new WorkflowConditionEvaluator(facts).Evaluate(expression);
}
catch (WorkflowConditionException ex)
{
throw new WorkflowEngineException(
$"MWORKFLOWCONFIG EVALCONDITION evaluation failed. " +
$"Expression: [{expression}]. " + ex.Message, ex);
}
}
///
/// Evaluates the EVALCONDITION of a MWORKFLOWDETAIL step against the facts.
/// Returns true when expression is null/empty (step always applies).
/// Throws a clear error when a field referenced in the expression is not in facts.
///
private bool EvaluateStepCondition(
string? expression,
IReadOnlyDictionary facts,
WorkflowStep step)
{
if (string.IsNullOrWhiteSpace(expression))
return true;
try
{
return new WorkflowConditionEvaluator(facts).Evaluate(expression);
}
catch (WorkflowConditionException ex)
{
throw new WorkflowEngineException(
$"MWORKFLOWDETAIL EVALCONDITION evaluation failed for " +
$"step '{step.DisplayName}' (Id={step.WorkflowDetailId}, " +
$"Level={step.ApprovalLevel}). " +
$"Expression: [{expression}]. " +
ex.Message, ex);
}
}
// ══════════════════════════════════════════════════════════════════════
// VALIDATION
// ══════════════════════════════════════════════════════════════════════
private static void ValidateCheckContext(WorkflowCheckContext ctx)
{
if (ctx == null)
throw new WorkflowEngineException(
"WorkflowCheckContext is null. " +
"Provide EntityId and Facts before calling CheckWorkFlowApplicability.");
if (ctx.EntityId == 0)
throw new WorkflowEngineException(
"WorkflowCheckContext.EntityId must be set to a valid entity id.");
}
private static void ValidateStartRequest(WorkflowStartRequest req)
{
if (req == null)
throw new WorkflowEngineException("WorkflowStartRequest is null.");
if (req.EntityId == 0)
throw new WorkflowEngineException(
"WorkflowStartRequest.EntityId must be set.");
// ObjectId may be 0 for WIP (form-based approval) — WipId is assigned inside the engine.
}
private static void ValidateWorkflowDefinition(WorkflowDefinition def)
{
if (def.Steps == null || !def.Steps.Any())
throw new WorkflowEngineException(
$"Workflow '{def.WorkflowName}' (Id={def.WorkflowId}) has no steps defined. " +
"Add at least one MWORKFLOWDETAIL row.");
var levels = def.GetApprovalLevels();
if (!levels.Any())
throw new WorkflowEngineException(
$"Workflow '{def.WorkflowName}' (Id={def.WorkflowId}) has no approval levels. " +
"Ensure MWORKFLOWDETAIL rows have a valid APPROVALLEVEL value.");
// Validate that every step has a valid assignment
foreach (var step in def.Steps)
{
// WorkflowAssignmentId is populated by the LEFT JOIN in GetWorkFlowSteps.
// A join miss (SQL NULL) maps to this non-nullable int's default, 0 — the
// exact sentinel for "no matching MWORKFLOWASSIGNMENT row". A real matched
// row's id is a negative GB5 autonumber value, so this must be == 0, never
// <= 0 (that would misfire on every correctly configured step).
// Note: MWORKFLOWDETAIL.ASSIGNMENTID is a FK, NOT a UserId/RoleId.
if (step.WorkflowAssignmentId == 0)
throw new WorkflowEngineException(
$"Step '{step.DisplayName}' (Id={step.WorkflowDetailId}, " +
$"Level={step.ApprovalLevel}) has ASSIGNMENTID={step.AssignmentId} " +
"but no matching row exists in MWORKFLOWASSIGNMENT. " +
"Insert a MWORKFLOWASSIGNMENT row with " +
$"WORKFLOWASSIGNMENTID={step.AssignmentId} and configure " +
"STRATEGYTYPE + RULEEXPRESSION.");
}
}
// ══════════════════════════════════════════════════════════════════════
// POST-COMPLETION ENTITY UPDATE
// ══════════════════════════════════════════════════════════════════════
///
/// Marks the instance as Completed and either:
/// - For WIP: marks TWORKFLOWWIP.STATUS=Approved and returns a WipDispatchInfo
/// so the caller can replay the original API call after committing.
/// - For regular: updates the entity record STATUS=1 directly.
///
private async Task<(WipDispatchInfo? Dispatch, WorkflowOutcomeInfo? Outcome)> FinaliseCompletionAsync(
WorkflowInstance instance,
LoginDTO login,
DbTransaction tx)
{
instance.WorkflowStatus = (int)WorkflowInstanceStatus.Completed;
await _runtime.UpdateInstanceAsync(instance, login, tx);
// Check if this is a WIP instance
var wip = await _runtime.GetWipByInstanceIdAsync(instance.WorkflowInstanceId, login, tx);
if (wip != null)
{
// WIP final approval: mark approved, return dispatch info.
// The WIP dispatcher replays the entity save with STATUS=1, which
// triggers the entity's own event — no separate outcome event needed.
await _runtime.UpdateWipStatusAsync(wip.WipId, 2, 1, 0, login.UserId, login, tx);
return (new WipDispatchInfo
{
WipId = wip.WipId,
WorkflowInstanceId = instance.WorkflowInstanceId,
ApiEndpoint = wip.ApiEndpoint,
DataJson = wip.DataJson,
LoginJson = wip.LoginJson
}, null);
}
// Regular workflow: update entity STATUS = 1, then emit approval outcome
// so MACTION subscribers can send approval notification emails.
await CompleteEntityStatusAsync(instance, login, tx);
return (null, new WorkflowOutcomeInfo
{
WorkflowInstanceId = instance.WorkflowInstanceId,
WorkflowId = instance.WorkflowId,
EntityId = instance.EntityId,
ObjectId = instance.ObjectId,
Action = (int)WorkflowAction.Approve,
ActionByUserId = login.UserId,
DataJson = instance.DataJson,
FactsJson = instance.FactsJson
});
}
///
/// After workflow completion, marks the source entity record as Approved (STATUS = 1).
///
/// Lookup chain:
/// TWORKFLOWINSTANCE.ENTITYID
/// → MENTITY.DBOBJECTID
/// → DBOBJECT.DBOBJECTNAME (e.g. 'TTASK')
/// → PK column of that table
/// → UPDATE {table} SET STATUS = 1 WHERE {pk} = OBJECTID
///
/// STATUS values: 0 = Pending, 1 = Approved.
/// Failures are swallowed — the workflow itself is already marked complete.
///
private Task CompleteEntityStatusAsync(
WorkflowInstance instance,
LoginDTO login,
DbTransaction tx)
=> SetEntityStatusAsync(instance, 1, login, tx);
private async Task SetEntityStatusAsync(
WorkflowInstance instance,
int status,
LoginDTO login,
DbTransaction tx)
{
try
{
// Step 1: resolve table name via MENTITY → DBOBJECT
string? tableName = await _qe.QuerySingleAsync(
login,
WorkFlowQB.GetDbObjectNameByEntityId,
new { entityid = instance.EntityId },
tx);
if (string.IsNullOrWhiteSpace(tableName) || !IsValidIdentifier(tableName))
return;
tableName = tableName.Trim();
// Step 2: resolve primary key column of the table
string? pkCol = await GetPrimaryKeyColumnAsync(tableName, login, tx);
if (string.IsNullOrWhiteSpace(pkCol))
return;
// Step 3: UPDATE {table} SET STATUS = {status} WHERE {pk} = OBJECTID
string updateSql =
$"UPDATE {tableName} SET STATUS = {status} WHERE {pkCol} = @objectid";
await _qe.ExecuteAsync(
login, updateSql, new { objectid = instance.ObjectId }, tx);
}
catch (Exception ex)
{
// Entity STATUS update is best-effort — never rolls back the workflow transaction.
_logger.LogError(ex,
"SetEntityStatusAsync: failed to update STATUS={Status} for EntityId={EntityId} ObjectId={ObjectId}. " +
"Tenant={TenantId} DB={DatabaseName}",
status, instance.EntityId, instance.ObjectId,
login.ClientId, login.DatabaseName);
}
}
private async Task GetPrimaryKeyColumnAsync(
string tableName,
LoginDTO login,
DbTransaction tx)
{
try
{
if (!IsValidIdentifier(tableName))
return null;
string sql = login.DatabaseType == DBType.SQL
? WorkFlowQB.GetPrimaryKeyColumnAsync
: WorkFlowQB.GetPrimaryKeyColumnAsync_pg;
return await _qe.QuerySingleAsync(
login, sql, new { tablename = tableName }, tx);
}
catch
{
return null;
}
}
private static bool IsValidIdentifier(string name)
=> !string.IsNullOrWhiteSpace(name) &&
name.All(c => char.IsLetterOrDigit(c) || c == '_');
// ══════════════════════════════════════════════════════════════════════
// HELPERS
// ══════════════════════════════════════════════════════════════════════
///
/// Converts a facts dictionary to an anonymous-object-style Dapper parameter bag.
/// Dapper accepts directly via DynamicParameters.
///
///
/// Rebuilds the facts dictionary from the two persisted JSON columns so that
/// every approval level can evaluate EVALCONDITION and resolve EnrichQualifier
/// (StrategyType=2) assignments without the caller re-supplying them.
///
/// Merge priority (lowest → highest):
/// 1. — enriched qualifier values saved at workflow start
/// 2. — entity DTO fields (override enriched on key clash)
/// 3. — caller-supplied action facts (highest priority)
///
/// This means a field like "ReportingToEmployeeId" written by a DB_ENRICH qualifier
/// is available at Level 2, 3, and every subsequent level without any additional
/// qualifier execution or caller-side re-enrichment.
///
private static IReadOnlyDictionary BuildFactsFromDataJson(
object? dataJson,
string? factsJson,
IReadOnlyDictionary overrides)
{
var facts = new Dictionary(StringComparer.OrdinalIgnoreCase);
// ── Layer 1: enriched qualifier facts (base) ──────────────────────
if (!string.IsNullOrWhiteSpace(factsJson))
{
try
{
var enriched = Newtonsoft.Json.JsonConvert.DeserializeObject<
Dictionary>(factsJson);
if (enriched != null)
foreach (var kv in enriched)
facts[kv.Key] = kv.Value;
}
catch
{
// Unparseable FactsJson — skip; DataJson and overrides still apply.
}
}
// ── Layer 2: entity DTO snapshot (overrides enriched on key clash) ─
if (dataJson != null)
{
try
{
string json = dataJson is string s
? s
: Newtonsoft.Json.JsonConvert.SerializeObject(dataJson);
// Guard against double-encoded JSON (legacy rows stored as quoted string).
if (!string.IsNullOrEmpty(json) && json.TrimStart().StartsWith('"'))
{
var inner = Newtonsoft.Json.JsonConvert.DeserializeObject(json);
if (!string.IsNullOrEmpty(inner))
json = inner;
}
var dict = Newtonsoft.Json.JsonConvert.DeserializeObject<
Dictionary>(json);
if (dict != null)
{
foreach (var kv in dict)
{
// Plain key — backward compat for any existing conditions
facts[kv.Key] = kv.Value;
// DTO. prefix — so EVALCONDITION can use "DTO.FieldName" syntax.
// Skip internal WIP marker fields (__WipIsNew, __WipPkField).
if (!kv.Key.StartsWith("__"))
facts[$"DTO.{kv.Key}"] = kv.Value;
}
}
}
catch
{
// Unparseable DataJson — proceed with enriched facts and overrides.
}
}
// ── Layer 3: caller-supplied action overrides (highest priority) ───
foreach (var kv in overrides)
facts[kv.Key] = kv.Value;
return facts;
}
///
/// Serialises the scalar entries of to a JSON string
/// suitable for storage in TWORKFLOWINSTANCE.FACTSJSON.
/// Non-scalar values (collections, nested DTOs) are excluded — they cannot be
/// meaningfully round-tripped and are not needed for EnrichQualifier resolution.
/// Returns null when the resulting map is empty.
///
private static string? SerializeFactsJson(IReadOnlyDictionary facts)
{
if (facts == null || facts.Count == 0)
return null;
var scalar = new Dictionary(StringComparer.OrdinalIgnoreCase);
foreach (var kv in facts)
{
if (kv.Value == null || IsScalarType(kv.Value.GetType()))
scalar[kv.Key] = kv.Value;
}
return scalar.Count > 0
? Newtonsoft.Json.JsonConvert.SerializeObject(scalar)
: null;
}
///
/// Parses (either a serialised JSON string or a live
/// object) and writes every scalar property into under
/// the "DTO." prefix — e.g. property "TaskDetailType" → key "DTO.TaskDetailType".
///
/// This lets MWORKFLOWCONFIG.EVALCONDITION and MWORKFLOWDETAIL.EVALCONDITION
/// reference any field of the incoming DTO dynamically without the caller having
/// to manually pre-populate the Facts dictionary.
///
/// Rules:
/// • Object and Array tokens are skipped (not scalar).
/// • Internal WIP marker fields (__WipIsNew, __WipPkField) are skipped.
/// • Numeric tokens are stored as so the evaluator's
/// numeric comparison path kicks in for == != > < >= <=.
/// • Existing keys in are overwritten so the DTO
/// snapshot always reflects the actual posted values.
///
private static void FlattenDataJsonFacts(object? dataJson, Dictionary target)
{
if (dataJson is null) return;
string jsonStr;
try
{
jsonStr = dataJson is string s ? s : JsonConvert.SerializeObject(dataJson);
// Guard against double-encoded JSON (legacy rows stored as quoted string).
if (!string.IsNullOrEmpty(jsonStr) && jsonStr.TrimStart().StartsWith('"'))
{
var inner = JsonConvert.DeserializeObject(jsonStr);
if (!string.IsNullOrEmpty(inner)) jsonStr = inner;
}
}
catch { return; }
// DataJson may be "{...}" (single DTO) or "[{...}]" (array endpoint, e.g. List).
// For array format, take the first element — WIP always stores one DTO per record.
JObject jObj;
try
{
var root = JToken.Parse(jsonStr);
if (root is JArray arr && arr.Count > 0 && arr[0] is JObject firstEl)
jObj = firstEl;
else if (root is JObject obj)
jObj = obj;
else
return; // empty array or unexpected shape — skip
}
catch { return; } // malformed JSON — skip silently
foreach (var prop in jObj.Properties())
{
if (prop.Name.StartsWith("__")) continue; // skip WIP internals
var token = prop.Value;
if (token.Type is JTokenType.Object or JTokenType.Array or JTokenType.None or JTokenType.Undefined)
continue;
object? scalar = token.Type switch
{
JTokenType.Integer => token.Value(),
JTokenType.Float => token.Value(),
JTokenType.Boolean => token.Value(),
JTokenType.Date => token.Value(),
JTokenType.Guid => token.Value(),
JTokenType.Null => null,
_ => token.Value()
};
target[$"DTO.{prop.Name}"] = scalar;
}
}
// SQL Server minimum date — DateTime.MinValue (0001-01-01) causes SqlDateTime overflow
private static readonly DateTime _sqlMinDate = new DateTime(1753, 1, 1);
private static object BuildSqlParameters(IReadOnlyDictionary facts)
{
var dp = new Dapper.DynamicParameters();
foreach (var kv in facts)
{
object? value = kv.Value;
// Skip null — nothing to bind
if (value is null)
{
dp.Add(kv.Key, null);
continue;
}
// Skip complex / collection types that Dapper cannot map to a SQL parameter.
// Facts are reflected from the full DTO (e.g. TLeaveDetailArray is a
// List) — only scalar primitives are valid SQL parameters.
var type = value.GetType();
if (!IsScalarType(type))
continue;
// Clamp out-of-range DateTime values to null before passing to SQL Server.
// Facts are reflected from the DTO — unset DateTime fields default to
// DateTime.MinValue (0001-01-01) which is below SQL Server's minimum (1753-01-01).
if (value is DateTime dt && dt < _sqlMinDate)
value = null;
dp.Add(kv.Key, value);
}
return dp;
}
///
/// Returns true for types that Dapper can safely bind as a SQL parameter:
/// primitives, strings, decimals, DateTimes, Guids, enums, and their
/// nullable counterparts. Collections, DTOs, and other complex objects
/// are excluded to prevent from Dapper.
///
private static bool IsScalarType(Type type)
{
// Unwrap Nullable
var underlying = Nullable.GetUnderlyingType(type) ?? type;
return underlying.IsPrimitive
|| underlying.IsEnum
|| underlying == typeof(string)
|| underlying == typeof(decimal)
|| underlying == typeof(DateTime)
|| underlying == typeof(DateTimeOffset)
|| underlying == typeof(Guid)
|| underlying == typeof(byte[]);
}
private static int GetDynamicInt(
IDictionary dict,
string key,
string stepName)
{
// Case-insensitive key lookup
var kv = dict.FirstOrDefault(
d => string.Equals(d.Key, key, StringComparison.OrdinalIgnoreCase));
if (kv.Value == null)
throw new WorkflowEngineException(
$"Dynamic SQL for step '{stepName}' returned a row where " +
$"'{key}' is NULL. The query must return a valid UserId.");
if (int.TryParse(kv.Value.ToString(), out int id))
return id;
throw new WorkflowEngineException(
$"Dynamic SQL for step '{stepName}' returned '{kv.Value}' " +
$"for column '{key}' which cannot be converted to an integer UserId.");
}
// ── Inner result type ─────────────────────────────────────────────────
private sealed class AssigneeResult
{
public int UserId { get; set; } = -1;
public int RoleId { get; set; } = -1;
public int GroupId { get; set; } = -1;
}
}
// ══════════════════════════════════════════════════════════════════════════
// DELEGATION RESOLVER
// ══════════════════════════════════════════════════════════════════════════
public class UserDelegationResolver : IUserDelegationResolver
{
private readonly IWorkFlowRunTime _runtime;
public UserDelegationResolver(IWorkFlowRunTime runtime)
{
_runtime = runtime;
}
public Task> GetAsync(
long userId, LoginDTO login, DbTransaction tx)
=> _runtime.GetUserDelegationsAsync(userId, login, tx);
}
// ══════════════════════════════════════════════════════════════════════════
// ENGINE EXCEPTION
// ══════════════════════════════════════════════════════════════════════════
///
/// Thrown for business-rule violations inside the workflow engine.
/// Always contains a clear, user-facing message suitable for display.
///
public sealed class WorkflowEngineException : Exception
{
public WorkflowEngineException(string message)
: base(message) { }
public WorkflowEngineException(string message, Exception inner)
: base(message, inner) { }
}
}