Files
HC900-Crawler/src/Infrastructure/Reporting/ReportMetricService.cs
windpacer 5f2780186e feat(report): P2-A 토큰 파서 일반화 — period=DAILY|MONTHLY|YEARLY 윈도 분기
- MetricRequest/Result에 Period 추가, ReportMetricService.PeriodWindowUtc로
  선택일이 속한 일/월/연 [from,to)를 KST→UTC 변환(DAILY는 기존 동작 동일).
- 오타 period는 daily 묵인 대신 명시적 error(결정론 게이트).
- ReportFillService 토큰 period= 파싱 + 셀주석/cells_json 박제.
- summary 엔드포인트 period 파라미터, 웹 바로보기 기간 드롭다운 + 치트시트.

Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com>
2026-06-17 21:07:03 +09:00

364 lines
18 KiB
C#

using System.Data;
using System.Data.Common;
using Hc900Crawler.Core.Application.DTOs;
using Hc900Crawler.Core.Application.Interfaces;
using Hc900Crawler.Infrastructure.Database;
using Microsoft.EntityFrameworkCore;
using Microsoft.Extensions.Logging;
namespace Hc900Crawler.Infrastructure.Reporting;
/// <summary>
/// 결정론 메트릭 엔진. history_table(60s)·fast_record(s/min) 동일 long 포맷이라 source만 교체.
/// raw SQL은 FeedforwardAuditService 패턴(_ctx.Database.GetDbConnection() + @param).
/// 표본 0 → no_data(0 날조 금지), 매핑 미정의 → error (결정론 검증 게이트).
/// </summary>
public sealed class ReportMetricService : IReportMetricService
{
private readonly Hc900DbContext _ctx;
private readonly ILogger<ReportMetricService> _logger;
private readonly ReportColumnMap _map;
private const string NUMERIC = "^-?[0-9]+(\\.[0-9]+)?$";
public ReportMetricService(Hc900DbContext ctx, ILogger<ReportMetricService> logger, ReportColumnMap map)
{ _ctx = ctx; _logger = logger; _map = map; }
private static readonly HashSet<string> QV_METRICS =
new() { "production_total", "yield_qv", "energy_intensity_qv", "mass_balance_closure" };
// recorded_at 윈도로 조회 가능한 시계열 소스(동일 long 포맷). history_1min_src=연속집계 호환뷰.
private static readonly HashSet<string> SERIES_SOURCES =
new() { "history_table", "history_1s", "history_1min_src" };
private static readonly HashSet<string> PERIODS =
new(StringComparer.OrdinalIgnoreCase) { "DAILY", "MONTHLY", "YEARLY" };
/// <summary>KST 날짜가 속한 일/월/연 경계 [from, to) 를 UTC로. recorded_at(UTC) 비교용.</summary>
private static (DateTime fromUtc, DateTime toUtc) PeriodWindowUtc(DateTime dateKst, string period)
{
var d = dateKst.Date;
var (fromKst, toKst) = period switch
{
"MONTHLY" => (new DateTime(d.Year, d.Month, 1), new DateTime(d.Year, d.Month, 1).AddMonths(1)),
"YEARLY" => (new DateTime(d.Year, 1, 1), new DateTime(d.Year, 1, 1).AddYears(1)),
_ => (d, d.AddDays(1)),
};
return (DateTime.SpecifyKind(fromKst, DateTimeKind.Unspecified).AddHours(-9),
DateTime.SpecifyKind(toKst, DateTimeKind.Unspecified).AddHours(-9));
}
public async Task<MetricResultDto> ComputeAsync(MetricRequestDto req, CancellationToken ct = default)
{
bool isFast = req.SourceTable == "fast_record";
// history_table(60s) | history_1s(1s 버퍼) | history_1min_src(연속집계) | fast_record. 미지정/미허용→history_table.
string tbl = isFast ? "fast_record"
: SERIES_SOURCES.Contains(req.SourceTable) ? req.SourceTable : "history_table";
string period = string.IsNullOrWhiteSpace(req.Period) ? "DAILY" : req.Period.Trim().ToUpperInvariant();
var res = new MetricResultDto
{
Metric = req.Metric, Column = req.Column, Period = period,
Source = tbl, SamplingMs = isFast ? 0 : (tbl == "history_1s" ? 1000 : 60000)
};
// 결정론 게이트: 오타 period를 daily로 묵인하면 틀린 값 → 명시적 에러.
if (!PERIODS.Contains(period))
{ res.Status = "error"; res.Error = $"미지원 period: {req.Period} (DAILY|MONTHLY|YEARLY)"; res.Value = null; return res; }
// PeriodDateKst가 속한 일/월/연 [from, to) KST → UTC (recorded_at은 UTC). fast_record는 session 기준이라 윈도 무관.
var (fromUtc, toUtc) = PeriodWindowUtc(req.PeriodDateKst, period);
try
{
var conn = _ctx.Database.GetDbConnection();
if (conn.State != ConnectionState.Open) await conn.OpenAsync(ct);
if (req.Metric == "dynamics")
{
await DynamicsAsync(res, conn, tbl, req.Column, fromUtc, toUtc, isFast, req.SessionId, ct);
}
else if (QV_METRICS.Contains(req.Metric))
{
if (!_map.TryResolveQv(req.Column, req.Metric, out var qspec) || qspec is null)
{ res.Status = "error"; res.Error = $"미정의 QV 매핑: {req.Column}/{req.Metric}"; res.Value = null; return res; }
res.Unit = qspec.Unit;
_map.TryResolveCleaning(req.Column, out var clean); // 비정상운전 제외 마스크(없으면 null)
await ComputeQvAsync(res, conn, tbl, qspec, clean, fromUtc, toUtc, isFast, req.SessionId, ct);
}
else
{
if (!_map.TryResolve(req.Column, req.Metric, out var spec) || spec is null)
{ res.Status = "error"; res.Error = $"미정의 매핑: {req.Column}/{req.Metric}"; res.Value = null; return res; }
res.Unit = spec.Unit;
if (req.Metric == "control_residual")
await ResidualAsync(res, conn, tbl, spec, fromUtc, toUtc, isFast, req.SessionId, ct);
else
await RatioAsync(res, conn, tbl, spec, fromUtc, toUtc, isFast, req.SessionId, ct);
}
if (isFast && res.Status == "ok") res.SamplingMs = await FastSamplingMsAsync(conn, req.SessionId ?? -1, ct);
}
catch (Exception ex)
{
_logger.LogError(ex, "[Report] metric 실패 {Col}/{M}", req.Column, req.Metric);
res.Status = "error"; res.Error = ex.Message; res.Value = null;
}
return res;
}
// ── 동특성(고해상 전용): 하부온도 루프 PV vs OP. 밸브 travel/hunting 진단. FOPDT 모델링 아님(step-test 필요). ──
private async Task DynamicsAsync(MetricResultDto res, DbConnection conn, string tbl, string column,
DateTime fromUtc, DateTime toUtc, bool isFast, int? sid, CancellationToken ct)
{
res.Unit = "OP/h";
if (tbl != "history_1s" && tbl != "fast_record")
{ res.Status = "no_data"; res.Value = null; res.Error = "동특성은 고해상 소스 필요(history_1s 또는 fast_record). 60초 history 불가."; return; }
if (!_map.TryResolveDynamics(column, out var pv, out var op) || pv is null || op is null)
{ res.Status = "error"; res.Value = null; res.Error = $"동특성 루프 미정의: {column}"; return; }
await using var cmd = conn.CreateCommand();
cmd.CommandText = $@"
WITH s AS (
SELECT date_trunc('second', recorded_at) ts,
max(CASE WHEN tagname=@pv THEN value::float END) pv,
max(CASE WHEN tagname=@op THEN value::float END) op
FROM hc900.{tbl}
WHERE tagname IN (@pv,@op) AND value ~ '{NUMERIC}'
AND ({(isFast ? "session_id = @sid" : "recorded_at >= @from AND recorded_at < @to")})
GROUP BY 1
), m AS (
SELECT avg(pv) pvm, stddev(pv) pvsd, count(*) n,
extract(epoch FROM (max(ts)-min(ts))) dur FROM s WHERE pv IS NOT NULL
), seq AS (
SELECT ts, sign(pv - (SELECT pvm FROM m)) side,
abs(op - lag(op) OVER (ORDER BY ts)) adop
FROM s WHERE pv IS NOT NULL
), seq2 AS (
SELECT side, lag(side) OVER (ORDER BY ts) ps, adop FROM seq
)
SELECT (SELECT pvsd FROM m), (SELECT n FROM m), (SELECT dur FROM m),
count(*) FILTER (WHERE side <> ps AND side <> 0 AND ps IS NOT NULL),
coalesce(sum(adop), 0)
FROM seq2";
AddP(cmd, "@pv", pv); AddP(cmd, "@op", op);
if (isFast) AddP(cmd, "@sid", sid ?? -1); else { AddP(cmd, "@from", fromUtc); AddP(cmd, "@to", toUtc); }
await using var rd = await cmd.ExecuteReaderAsync(ct);
if (await rd.ReadAsync(ct) && !rd.IsDBNull(1) && rd.GetInt64(1) >= 30)
{
double pvsd = rd.IsDBNull(0) ? 0 : rd.GetDouble(0);
res.N = (int)rd.GetInt64(1);
double dur = rd.IsDBNull(2) ? 0 : rd.GetDouble(2);
long crossings = rd.IsDBNull(3) ? 0 : rd.GetInt64(3);
double opTravel = rd.IsDBNull(4) ? 0 : rd.GetDouble(4);
res.Value = dur > 0 ? opTravel * 3600.0 / dur : null; // OP travel/h = 밸브 활동(stiction/hunting/마모 proxy)
res.Extra["pv_sd"] = pvsd; // PV 변동(고해상)
res.Extra["osc_period_s"] = crossings > 0 && dur > 0 ? 2.0 * dur / crossings : null; // 헌팅 주기(s)
res.Extra["crossings"] = crossings;
res.Extra["dur_s"] = dur;
if (res.Value is null) res.Status = "no_data";
}
else { res.Status = "no_data"; res.Value = null; res.Error = "고해상 표본 부족(≥30s 필요)"; }
}
// ── 적산(.QV) 메트릭: Single=ΔA, Ratio=ΔA/ΔB, Closure=100·ΣΔOut/Δfeed. cleaning/drawdown 제외. ──
private async Task ComputeQvAsync(MetricResultDto res, DbConnection conn, string tbl, QvSpec s,
CleaningSpec? cl, DateTime fromUtc, DateTime toUtc, bool isFast, int? sid, CancellationToken ct)
{
if (s.Kind == QvKind.Single)
{
var (tot, n, excl) = await QvDeltaAsync(conn, tbl, s.A, cl, fromUtc, toUtc, isFast, sid, ct);
if (tot is null) { res.Status = "no_data"; res.Value = null; return; }
res.Value = tot; res.N = n; res.Extra["excluded_min"] = excl;
}
else if (s.Kind == QvKind.Ratio)
{
var (a, na, ea) = await QvDeltaAsync(conn, tbl, s.A, cl, fromUtc, toUtc, isFast, sid, ct);
var (b, nb, eb) = await QvDeltaAsync(conn, tbl, s.B!, cl, fromUtc, toUtc, isFast, sid, ct);
if (a is null || b is null || b == 0) { res.Status = "no_data"; res.Value = null; return; }
res.Value = a / b; res.N = Math.Min(na, nb);
res.Extra["numer_qv"] = a; res.Extra["denom_qv"] = b; res.Extra["excluded_min"] = ea;
// 물리적 타당성 게이트: 비율이 음수/0이하 또는 상한 초과 → 데이터 이상
if (res.Value <= 0 || (s.Max is double mx && res.Value > mx))
{
res.Status = "no_data";
res.Error = $"비물리적 비율 {res.Value:F1} (분모 적산 Δ={b:F0} 비정상 추정)";
res.Value = null;
}
}
else // Closure
{
var (feed, nf, ef) = await QvDeltaAsync(conn, tbl, s.A, cl, fromUtc, toUtc, isFast, sid, ct);
if (feed is null || feed == 0) { res.Status = "no_data"; res.Value = null; return; }
double outSum = 0;
for (int i = 0; i < s.Outputs.Count; i++)
{
var (d, _, _) = await QvDeltaAsync(conn, tbl, s.Outputs[i], cl, fromUtc, toUtc, isFast, sid, ct);
outSum += d ?? 0;
res.Extra[$"out{i}_qv"] = d; // out0=제품, out1=경비물, out2=중비물 (config 순서)
}
res.Value = 100.0 * outSum / feed; // 폐합 %
res.N = nf;
res.Extra["feed_qv"] = feed;
res.Extra["out_total"] = outSum;
res.Extra["product_qv"] = s.Outputs.Count > 0 ? res.Extra["out0_qv"] : null;
res.Extra["excluded_min"] = ef;
}
}
/// <summary>
/// 적산 Δ. 비정상운전 분(cleaning=진공高/제품~0, drawdown=feed~0) 제외 후, 정상구간 양(+)증분만 합산(cap 5e4).
/// 리셋 3종 자동처리: 999999 wrap·cleaning 리셋·운전조건변경 리셋. cl=null이면 마스크 없이 양증분합산.
/// 반환: (Δ합, 정상분수, 제외분수).
/// </summary>
private static async Task<(double? Total, int NNormal, int NExcluded)> QvDeltaAsync(
DbConnection conn, string tbl, string tag, CleaningSpec? cl,
DateTime fromUtc, DateTime toUtc, bool isFast, int? sid, CancellationToken ct)
{
bool hasVac = cl != null && !string.IsNullOrEmpty(cl.VacTag);
var pivot = new System.Text.StringBuilder("max(CASE WHEN tagname=@tag THEN value::float END) v");
var inTags = new List<string> { "@tag" };
var mask = new List<string>();
if (cl != null)
{
pivot.Append(", max(CASE WHEN tagname=@ptag THEN value::float END) prod");
pivot.Append(", max(CASE WHEN tagname=@ftag THEN value::float END) feed");
mask.Add("prod < @pmin OR prod IS NULL");
mask.Add("feed < @fmin OR feed IS NULL");
inTags.Add("@ptag"); inTags.Add("@ftag");
if (hasVac)
{
pivot.Append(", max(CASE WHEN tagname=@vtag THEN value::float END) vac");
mask.Add("vac > @vmax OR vac IS NULL");
inTags.Add("@vtag");
}
}
string cleanExpr = mask.Count > 0 ? "(" + string.Join(" OR ", mask) + ")" : "false";
string window = isFast ? "session_id = @sid" : "recorded_at >= @from AND recorded_at < @to";
await using var cmd = conn.CreateCommand();
cmd.CommandText = $@"
WITH pm AS (
SELECT date_trunc('minute', recorded_at) ts, {pivot}
FROM hc900.{tbl}
WHERE tagname IN ({string.Join(",", inTags)}) AND value ~ '{NUMERIC}' AND ({window})
GROUP BY 1
), fl AS (
SELECT ts, v, ({cleanExpr}) AS clean FROM pm WHERE v IS NOT NULL
), seq AS (
SELECT ts, v, clean,
lag(v) OVER (ORDER BY ts) pv,
lag(clean) OVER (ORDER BY ts) pc
FROM fl
)
SELECT coalesce(sum(v - pv) FILTER (
WHERE pv IS NOT NULL AND NOT clean AND NOT coalesce(pc, true)
AND v >= pv AND v - pv < 50000), 0),
count(*) FILTER (WHERE NOT clean),
count(*) FILTER (WHERE clean)
FROM seq;";
AddP(cmd, "@tag", tag);
if (cl != null)
{
AddP(cmd, "@ptag", cl.ProductTag); AddP(cmd, "@ftag", cl.FeedTag);
AddP(cmd, "@pmin", cl.ProductMin); AddP(cmd, "@fmin", cl.FeedMin);
if (hasVac) { AddP(cmd, "@vtag", cl.VacTag!); AddP(cmd, "@vmax", cl.VacMax); }
}
if (isFast) AddP(cmd, "@sid", sid ?? -1); else { AddP(cmd, "@from", fromUtc); AddP(cmd, "@to", toUtc); }
await using var rd = await cmd.ExecuteReaderAsync(ct);
if (await rd.ReadAsync(ct) && !rd.IsDBNull(0))
return (rd.GetDouble(0), (int)rd.GetInt64(1), (int)rd.GetInt64(2));
return (null, 0, 0);
}
// ── 비율: median(A/B) (분 버킷 피벗, 둘 다 유효한 분만). 클린율 = 1 - 2*good/raw ──
private async Task RatioAsync(MetricResultDto res, DbConnection conn, string tbl, MetricSpec s,
DateTime fromUtc, DateTime toUtc, bool isFast, int? sid, CancellationToken ct)
{
await using var cmd = conn.CreateCommand();
cmd.CommandText = $@"
WITH raw AS (
SELECT date_trunc('minute', recorded_at) ts, tagname, value::float v
FROM hc900.{tbl}
WHERE tagname IN (@a,@b) AND value ~ '{NUMERIC}'
AND ({(isFast ? "session_id = @sid" : "recorded_at >= @from AND recorded_at < @to")})
), tot AS (SELECT count(*) c FROM raw),
piv AS (
SELECT ts, max(v) FILTER (WHERE tagname=@a) a, max(v) FILTER (WHERE tagname=@b) b
FROM raw GROUP BY ts
), good AS (
SELECT a, b FROM piv
WHERE a BETWEEN @alo AND @ahi AND b BETWEEN @blo AND @bhi AND b <> 0
)
SELECT (SELECT percentile_cont(0.5) WITHIN GROUP (ORDER BY a/b) FROM good),
(SELECT count(*) FROM good),
(SELECT c FROM tot);";
AddP(cmd, "@a", s.A.Tag); AddP(cmd, "@b", s.B.Tag);
AddP(cmd, "@alo", s.A.Lo); AddP(cmd, "@ahi", s.A.Hi);
AddP(cmd, "@blo", s.B.Lo); AddP(cmd, "@bhi", s.B.Hi);
if (isFast) AddP(cmd, "@sid", sid ?? -1); else { AddP(cmd, "@from", fromUtc); AddP(cmd, "@to", toUtc); }
await using var rd = await cmd.ExecuteReaderAsync(ct);
if (await rd.ReadAsync(ct) && !rd.IsDBNull(0))
{
res.Value = rd.GetDouble(0);
res.N = (int)rd.GetInt64(1);
var nRaw = rd.GetInt64(2);
res.CleanedFraction = nRaw > 0 ? Math.Clamp(1.0 - (2.0 * res.N) / nRaw, 0, 1) : 0;
}
else { res.Status = "no_data"; res.Value = null; }
}
// ── 제어잔차: (PV-SP) mean/sd/abs_p95/이탈% ──
private async Task ResidualAsync(MetricResultDto res, DbConnection conn, string tbl, MetricSpec s,
DateTime fromUtc, DateTime toUtc, bool isFast, int? sid, CancellationToken ct)
{
await using var cmd = conn.CreateCommand();
cmd.CommandText = $@"
WITH raw AS (
SELECT date_trunc('minute', recorded_at) ts, tagname, value::float v
FROM hc900.{tbl}
WHERE tagname IN (@pv,@sp) AND value ~ '{NUMERIC}'
AND ({(isFast ? "session_id = @sid" : "recorded_at >= @from AND recorded_at < @to")})
), piv AS (
SELECT ts, max(v) FILTER (WHERE tagname=@pv) pv, max(v) FILTER (WHERE tagname=@sp) sp
FROM raw GROUP BY ts
), good AS (
SELECT pv - sp AS e FROM piv
WHERE pv BETWEEN @plo AND @phi AND sp BETWEEN @slo AND @shi
)
SELECT avg(e), stddev(e),
percentile_cont(0.95) WITHIN GROUP (ORDER BY abs(e)),
count(*),
100.0 * count(*) FILTER (WHERE abs(e) > 0.5) / NULLIF(count(*),0)
FROM good;";
AddP(cmd, "@pv", s.A.Tag); AddP(cmd, "@sp", s.B.Tag);
AddP(cmd, "@plo", s.A.Lo); AddP(cmd, "@phi", s.A.Hi);
AddP(cmd, "@slo", s.B.Lo); AddP(cmd, "@shi", s.B.Hi);
if (isFast) AddP(cmd, "@sid", sid ?? -1); else { AddP(cmd, "@from", fromUtc); AddP(cmd, "@to", toUtc); }
await using var rd = await cmd.ExecuteReaderAsync(ct);
if (await rd.ReadAsync(ct) && !rd.IsDBNull(3) && rd.GetInt64(3) > 0)
{
res.Value = rd.IsDBNull(0) ? null : rd.GetDouble(0);
res.Extra["sd"] = rd.IsDBNull(1) ? null : rd.GetDouble(1);
res.Extra["abs_p95"] = rd.IsDBNull(2) ? null : rd.GetDouble(2);
res.N = (int)rd.GetInt64(3);
res.Extra["out_pct_0_5"] = rd.IsDBNull(4) ? null : rd.GetDouble(4);
}
else { res.Status = "no_data"; res.Value = null; }
}
private static async Task<int> FastSamplingMsAsync(DbConnection conn, int sid, CancellationToken ct)
{
await using var c = conn.CreateCommand();
c.CommandText = "SELECT sampling_ms FROM hc900.fast_session WHERE id=@id";
AddP(c, "@id", sid);
var o = await c.ExecuteScalarAsync(ct);
return o is null or DBNull ? 0 : Convert.ToInt32(o);
}
private static void AddP(DbCommand cmd, string name, object val)
{ var p = cmd.CreateParameter(); p.ParameterName = name; p.Value = val ?? DBNull.Value; cmd.Parameters.Add(p); }
}