diff --git a/Assets/Dashboards/MergerLogsDashboard.json b/Assets/Dashboards/MergerLogsDashboard.json new file mode 100644 index 00000000..c0dd0dff --- /dev/null +++ b/Assets/Dashboards/MergerLogsDashboard.json @@ -0,0 +1,1134 @@ +{ + "title": "GPKG-Merger — Logs", + "uid": "gpkg-merger-logs", + "tags": [ + "gpkg-merger", + "logs" + ], + "timezone": "", + "schemaVersion": 39, + "version": 1, + "editable": true, + "refresh": "30s", + "time": { + "from": "now-3h", + "to": "now" + }, + "templating": { + "list": [ + { + "name": "loki", + "label": "Loki datasource", + "type": "datasource", + "query": "loki", + "current": {}, + "hide": 0, + "includeAll": false, + "multi": false, + "refresh": 1, + "regex": "", + "skipUrlSync": false + }, + { + "name": "namespace", + "label": "namespace", + "type": "query", + "datasource": { + "type": "loki", + "uid": "${loki}" + }, + "query": { + "label": "k8s_namespace_name", + "stream": "{k8s_deployment_name=~\"gpkg-merger.*\"}", + "type": 4, + "refId": "LokiVariableQueryEditor-VariableQuery" + }, + "definition": "label_values({k8s_deployment_name=~\"gpkg-merger.*\"}, k8s_namespace_name)", + "current": { + "text": "All", + "value": "$__all", + "selected": true + }, + "includeAll": true, + "multi": true, + "refresh": 2, + "regex": "", + "sort": 1, + "hide": 0, + "skipUrlSync": false + }, + { + "name": "job", + "label": "jobId", + "type": "textbox", + "query": "", + "current": { + "text": "", + "value": "" + }, + "hide": 0, + "skipUrlSync": false + }, + { + "name": "task", + "label": "taskId", + "type": "textbox", + "query": "", + "current": { + "text": "", + "value": "" + }, + "hide": 0, + "skipUrlSync": false + } + ] + }, + "panels": [ + { + "id": 1, + "type": "stat", + "title": "Added", + "datasource": { + "type": "loki", + "uid": "${loki}" + }, + "gridPos": { + "h": 5, + "w": 4, + "x": 0, + "y": 0 + }, + "fieldConfig": { + "defaults": { + "unit": "short", + "noValue": "No data", + "color": { + "mode": "thresholds" + }, + "thresholds": { + "mode": "absolute", + "steps": [ + { + "color": "blue", + "value": null + } + ] + } + }, + "overrides": [] + }, + "options": { + "reduceOptions": { + "calcs": [ + "lastNotNull" + ], + "fields": "", + "values": false + }, + "colorMode": "value", + "graphMode": "area" + }, + "targets": [ + { + "refId": "A", + "datasource": { + "type": "loki", + "uid": "${loki}" + }, + "editorMode": "code", + "expr": "sum(last_over_time({k8s_namespace_name=~\"$namespace\", k8s_deployment_name=~\"gpkg-merger.*\"} | json | __error__=\"\" | jobId=~\"$job\" | taskId=~\"$task\" | message=~\".*Merge report.*\" | unwrap added [$__range]))", + "queryType": "instant" + } + ] + }, + { + "id": 2, + "type": "stat", + "title": "Merged", + "datasource": { + "type": "loki", + "uid": "${loki}" + }, + "gridPos": { + "h": 5, + "w": 4, + "x": 4, + "y": 0 + }, + "fieldConfig": { + "defaults": { + "unit": "short", + "noValue": "No data", + "color": { + "mode": "thresholds" + }, + "thresholds": { + "mode": "absolute", + "steps": [ + { + "color": "blue", + "value": null + } + ] + } + }, + "overrides": [] + }, + "options": { + "reduceOptions": { + "calcs": [ + "lastNotNull" + ], + "fields": "", + "values": false + }, + "colorMode": "value", + "graphMode": "area" + }, + "targets": [ + { + "refId": "A", + "datasource": { + "type": "loki", + "uid": "${loki}" + }, + "editorMode": "code", + "expr": "sum(last_over_time({k8s_namespace_name=~\"$namespace\", k8s_deployment_name=~\"gpkg-merger.*\"} | json | __error__=\"\" | jobId=~\"$job\" | taskId=~\"$task\" | message=~\".*Merge report.*\" | unwrap merged [$__range]))", + "queryType": "instant" + } + ] + }, + { + "id": 3, + "type": "stat", + "title": "Replaced", + "datasource": { + "type": "loki", + "uid": "${loki}" + }, + "gridPos": { + "h": 5, + "w": 4, + "x": 8, + "y": 0 + }, + "fieldConfig": { + "defaults": { + "unit": "short", + "noValue": "No data", + "color": { + "mode": "thresholds" + }, + "thresholds": { + "mode": "absolute", + "steps": [ + { + "color": "blue", + "value": null + } + ] + } + }, + "overrides": [] + }, + "options": { + "reduceOptions": { + "calcs": [ + "lastNotNull" + ], + "fields": "", + "values": false + }, + "colorMode": "value", + "graphMode": "area" + }, + "targets": [ + { + "refId": "A", + "datasource": { + "type": "loki", + "uid": "${loki}" + }, + "editorMode": "code", + "expr": "sum(last_over_time({k8s_namespace_name=~\"$namespace\", k8s_deployment_name=~\"gpkg-merger.*\"} | json | __error__=\"\" | jobId=~\"$job\" | taskId=~\"$task\" | message=~\".*Merge report.*\" | unwrap replaced [$__range]))", + "queryType": "instant" + } + ] + }, + { + "id": 4, + "type": "stat", + "title": "Skipped", + "datasource": { + "type": "loki", + "uid": "${loki}" + }, + "gridPos": { + "h": 5, + "w": 4, + "x": 12, + "y": 0 + }, + "fieldConfig": { + "defaults": { + "unit": "short", + "noValue": "No data", + "color": { + "mode": "thresholds" + }, + "thresholds": { + "mode": "absolute", + "steps": [ + { + "color": "blue", + "value": null + } + ] + } + }, + "overrides": [] + }, + "options": { + "reduceOptions": { + "calcs": [ + "lastNotNull" + ], + "fields": "", + "values": false + }, + "colorMode": "value", + "graphMode": "area" + }, + "targets": [ + { + "refId": "A", + "datasource": { + "type": "loki", + "uid": "${loki}" + }, + "editorMode": "code", + "expr": "sum(last_over_time({k8s_namespace_name=~\"$namespace\", k8s_deployment_name=~\"gpkg-merger.*\"} | json | __error__=\"\" | jobId=~\"$job\" | taskId=~\"$task\" | message=~\".*Merge report.*\" | unwrap skipped [$__range]))", + "queryType": "instant" + } + ] + }, + { + "id": 5, + "type": "stat", + "title": "Total tiles", + "datasource": { + "type": "loki", + "uid": "${loki}" + }, + "gridPos": { + "h": 5, + "w": 4, + "x": 16, + "y": 0 + }, + "fieldConfig": { + "defaults": { + "unit": "short", + "noValue": "No data", + "color": { + "mode": "thresholds" + }, + "thresholds": { + "mode": "absolute", + "steps": [ + { + "color": "blue", + "value": null + } + ] + } + }, + "overrides": [] + }, + "options": { + "reduceOptions": { + "calcs": [ + "lastNotNull" + ], + "fields": "", + "values": false + }, + "colorMode": "value", + "graphMode": "area" + }, + "targets": [ + { + "refId": "A", + "datasource": { + "type": "loki", + "uid": "${loki}" + }, + "editorMode": "code", + "expr": "sum(last_over_time({k8s_namespace_name=~\"$namespace\", k8s_deployment_name=~\"gpkg-merger.*\"} | json | __error__=\"\" | jobId=~\"$job\" | taskId=~\"$task\" | message=~\".*Merge report.*\" | unwrap total [$__range]))", + "queryType": "instant" + } + ] + }, + { + "id": 6, + "type": "stat", + "title": "p95 duration", + "datasource": { + "type": "loki", + "uid": "${loki}" + }, + "gridPos": { + "h": 5, + "w": 4, + "x": 20, + "y": 0 + }, + "fieldConfig": { + "defaults": { + "unit": "s", + "noValue": "No data", + "color": { + "mode": "thresholds" + }, + "thresholds": { + "mode": "absolute", + "steps": [ + { + "color": "green", + "value": null + }, + { + "color": "orange", + "value": 60 + }, + { + "color": "red", + "value": 120 + } + ] + } + }, + "overrides": [] + }, + "options": { + "reduceOptions": { + "calcs": [ + "lastNotNull" + ], + "fields": "", + "values": false + }, + "colorMode": "value", + "graphMode": "area" + }, + "targets": [ + { + "refId": "A", + "datasource": { + "type": "loki", + "uid": "${loki}" + }, + "editorMode": "code", + "expr": "max(quantile_over_time(0.95, {k8s_namespace_name=~\"$namespace\", k8s_deployment_name=~\"gpkg-merger.*\"} | json | __error__=\"\" | jobId=~\"$job\" | taskId=~\"$task\" | message=~\".*Merge report.*\" | unwrap durationSeconds [$__range]))", + "queryType": "instant" + } + ] + }, + { + "id": 7, + "type": "timeseries", + "title": "Tiles by outcome (per $__interval)", + "datasource": { + "type": "loki", + "uid": "${loki}" + }, + "gridPos": { + "h": 8, + "w": 12, + "x": 0, + "y": 23 + }, + "fieldConfig": { + "defaults": { + "custom": { + "drawStyle": "bars", + "fillOpacity": 50, + "stacking": { + "mode": "normal" + } + } + }, + "overrides": [] + }, + "options": { + "legend": { + "displayMode": "list", + "placement": "bottom", + "showLegend": true + } + }, + "targets": [ + { + "refId": "A", + "datasource": { + "type": "loki", + "uid": "${loki}" + }, + "editorMode": "code", + "expr": "sum(sum_over_time({k8s_namespace_name=~\"$namespace\", k8s_deployment_name=~\"gpkg-merger.*\"} | json | __error__=\"\" | jobId=~\"$job\" | taskId=~\"$task\" | message=~\".*Merge report.*\" | unwrap added [$__interval]))", + "queryType": "range", + "legendFormat": "added" + }, + { + "refId": "B", + "datasource": { + "type": "loki", + "uid": "${loki}" + }, + "editorMode": "code", + "expr": "sum(sum_over_time({k8s_namespace_name=~\"$namespace\", k8s_deployment_name=~\"gpkg-merger.*\"} | json | __error__=\"\" | jobId=~\"$job\" | taskId=~\"$task\" | message=~\".*Merge report.*\" | unwrap merged [$__interval]))", + "queryType": "range", + "legendFormat": "merged" + }, + { + "refId": "C", + "datasource": { + "type": "loki", + "uid": "${loki}" + }, + "editorMode": "code", + "expr": "sum(sum_over_time({k8s_namespace_name=~\"$namespace\", k8s_deployment_name=~\"gpkg-merger.*\"} | json | __error__=\"\" | jobId=~\"$job\" | taskId=~\"$task\" | message=~\".*Merge report.*\" | unwrap replaced [$__interval]))", + "queryType": "range", + "legendFormat": "replaced" + }, + { + "refId": "D", + "datasource": { + "type": "loki", + "uid": "${loki}" + }, + "editorMode": "code", + "expr": "sum(sum_over_time({k8s_namespace_name=~\"$namespace\", k8s_deployment_name=~\"gpkg-merger.*\"} | json | __error__=\"\" | jobId=~\"$job\" | taskId=~\"$task\" | message=~\".*Merge report.*\" | unwrap skipped [$__interval]))", + "queryType": "range", + "legendFormat": "skipped" + } + ] + }, + { + "id": 8, + "type": "timeseries", + "title": "Task duration (s)", + "datasource": { + "type": "loki", + "uid": "${loki}" + }, + "gridPos": { + "h": 8, + "w": 12, + "x": 12, + "y": 23 + }, + "fieldConfig": { + "defaults": { + "unit": "s", + "custom": { + "drawStyle": "line", + "fillOpacity": 10 + } + }, + "overrides": [] + }, + "options": { + "legend": { + "displayMode": "list", + "placement": "bottom", + "showLegend": true + } + }, + "targets": [ + { + "refId": "A", + "datasource": { + "type": "loki", + "uid": "${loki}" + }, + "editorMode": "code", + "expr": "max(quantile_over_time(0.5, {k8s_namespace_name=~\"$namespace\", k8s_deployment_name=~\"gpkg-merger.*\"} | json | __error__=\"\" | jobId=~\"$job\" | taskId=~\"$task\" | message=~\".*Merge report.*\" | unwrap durationSeconds [$__interval]))", + "queryType": "range", + "legendFormat": "p50" + }, + { + "refId": "B", + "datasource": { + "type": "loki", + "uid": "${loki}" + }, + "editorMode": "code", + "expr": "max(quantile_over_time(0.95, {k8s_namespace_name=~\"$namespace\", k8s_deployment_name=~\"gpkg-merger.*\"} | json | __error__=\"\" | jobId=~\"$job\" | taskId=~\"$task\" | message=~\".*Merge report.*\" | unwrap durationSeconds [$__interval]))", + "queryType": "range", + "legendFormat": "p95" + } + ] + }, + { + "id": 9, + "type": "table", + "title": "Per-task merge reports (click job/task to filter)", + "datasource": { + "type": "loki", + "uid": "${loki}" + }, + "gridPos": { + "h": 10, + "w": 24, + "x": 0, + "y": 13 + }, + "fieldConfig": { + "defaults": { + "custom": { + "filterable": true, + "align": "auto" + } + }, + "overrides": [ + { + "matcher": { + "id": "byName", + "options": "job" + }, + "properties": [ + { + "id": "links", + "value": [ + { + "title": "Filter dashboard to this job", + "url": "/d/gpkg-merger-logs?${__url_time_range}&var-loki=${loki}&${namespace:queryparam}&var-job=${__data.fields.job}&var-task=", + "targetBlank": false + } + ] + } + ] + }, + { + "matcher": { + "id": "byName", + "options": "task" + }, + "properties": [ + { + "id": "links", + "value": [ + { + "title": "Filter dashboard to this job and task", + "url": "/d/gpkg-merger-logs?${__url_time_range}&var-loki=${loki}&${namespace:queryparam}&var-job=${__data.fields.job}&var-task=${__data.fields.task}", + "targetBlank": false + } + ] + } + ] + } + ] + }, + "options": { + "showHeader": true, + "cellHeight": "sm", + "footer": { + "show": false + } + }, + "targets": [ + { + "refId": "A", + "datasource": { + "type": "loki", + "uid": "${loki}" + }, + "editorMode": "code", + "expr": "{k8s_namespace_name=~\"$namespace\", k8s_deployment_name=~\"gpkg-merger.*\"} | json | __error__=\"\" | jobId=~\"$job\" | taskId=~\"$task\" | message=~\".*Merge report.*\"", + "queryType": "range" + } + ], + "transformations": [ + { + "id": "extractFields", + "options": { + "source": "labels", + "replace": false, + "keepTime": false + } + }, + { + "id": "organize", + "options": { + "excludeByName": { + "service": true, + "version": true, + "thread": true, + "service_name": true, + "k8s_container_name": true, + "k8s_deployment_name": true, + "k8s_namespace_name": true, + "k8s_pod_name": true, + "time": true, + "Line": true, + "labels": true, + "id": true, + "tsNs": true, + "level": true, + "category": true, + "message": true + }, + "indexByName": { + "Time": 0, + "jobId": 1, + "taskId": 2, + "added": 3, + "merged": 4, + "replaced": 5, + "skipped": 6, + "total": 7, + "durationSeconds": 8 + }, + "renameByName": { + "jobId": "job", + "taskId": "task", + "durationSeconds": "duration(s)" + } + } + } + ] + }, + { + "id": 10, + "type": "logs", + "title": "Task logs", + "datasource": { + "type": "loki", + "uid": "${loki}" + }, + "gridPos": { + "h": 11, + "w": 16, + "x": 0, + "y": 31 + }, + "options": { + "showTime": true, + "showLabels": false, + "showCommonLabels": false, + "wrapLogMessage": true, + "prettifyLogMessage": false, + "enableLogDetails": true, + "dedupStrategy": "none", + "sortOrder": "Descending" + }, + "targets": [ + { + "refId": "A", + "datasource": { + "type": "loki", + "uid": "${loki}" + }, + "editorMode": "code", + "expr": "{k8s_namespace_name=~\"$namespace\", k8s_deployment_name=~\"gpkg-merger.*\"} | json | __error__=\"\" | jobId=~\"$job\" | taskId=~\"$task\"", + "queryType": "range" + } + ] + }, + { + "id": 11, + "type": "timeseries", + "title": "Log rate by level", + "datasource": { + "type": "loki", + "uid": "${loki}" + }, + "gridPos": { + "h": 6, + "w": 8, + "x": 16, + "y": 31 + }, + "fieldConfig": { + "defaults": { + "custom": { + "drawStyle": "bars", + "fillOpacity": 40, + "stacking": { + "mode": "normal" + } + } + }, + "overrides": [] + }, + "options": { + "legend": { + "displayMode": "list", + "placement": "bottom", + "showLegend": true + } + }, + "targets": [ + { + "refId": "A", + "datasource": { + "type": "loki", + "uid": "${loki}" + }, + "editorMode": "code", + "expr": "sum by (level) (count_over_time({k8s_namespace_name=~\"$namespace\", k8s_deployment_name=~\"gpkg-merger.*\"} | json | __error__=\"\" | jobId=~\"$job\" | taskId=~\"$task\" [$__interval]))", + "queryType": "range", + "legendFormat": "{{level}}" + } + ] + }, + { + "id": 12, + "type": "stat", + "title": "Log lines (range)", + "datasource": { + "type": "loki", + "uid": "${loki}" + }, + "gridPos": { + "h": 5, + "w": 8, + "x": 16, + "y": 37 + }, + "fieldConfig": { + "defaults": { + "color": { + "mode": "thresholds" + }, + "thresholds": { + "mode": "absolute", + "steps": [ + { + "color": "green", + "value": null + } + ] + } + }, + "overrides": [] + }, + "options": { + "reduceOptions": { + "calcs": [ + "lastNotNull" + ], + "fields": "", + "values": false + }, + "colorMode": "value", + "graphMode": "area" + }, + "targets": [ + { + "refId": "A", + "datasource": { + "type": "loki", + "uid": "${loki}" + }, + "editorMode": "code", + "expr": "sum(count_over_time({k8s_namespace_name=~\"$namespace\", k8s_deployment_name=~\"gpkg-merger.*\"} | json | __error__=\"\" | jobId=~\"$job\" | taskId=~\"$task\" [$__range]))", + "queryType": "instant" + } + ] + }, + { + "id": 15, + "type": "stat", + "title": "Errors (range)", + "datasource": { + "type": "loki", + "uid": "${loki}" + }, + "gridPos": { + "h": 8, + "w": 4, + "x": 0, + "y": 5 + }, + "fieldConfig": { + "defaults": { + "unit": "short", + "color": { + "mode": "thresholds" + }, + "thresholds": { + "mode": "absolute", + "steps": [ + { + "color": "green", + "value": null + }, + { + "color": "red", + "value": 1 + } + ] + } + }, + "overrides": [] + }, + "options": { + "reduceOptions": { + "calcs": [ + "lastNotNull" + ], + "fields": "", + "values": false + }, + "colorMode": "background", + "graphMode": "area" + }, + "targets": [ + { + "refId": "A", + "datasource": { + "type": "loki", + "uid": "${loki}" + }, + "editorMode": "code", + "expr": "sum(count_over_time({k8s_namespace_name=~\"$namespace\", k8s_deployment_name=~\"gpkg-merger.*\"} | json | __error__=\"\" | jobId=~\"$job\" | taskId=~\"$task\" | level=~\"Error|Critical\" [$__range]))", + "queryType": "instant" + } + ] + }, + { + "id": 13, + "type": "table", + "title": "Error messages", + "datasource": { + "type": "loki", + "uid": "${loki}" + }, + "gridPos": { + "h": 8, + "w": 20, + "x": 4, + "y": 5 + }, + "fieldConfig": { + "defaults": { + "custom": { + "filterable": true, + "align": "auto" + } + }, + "overrides": [ + { + "matcher": { + "id": "byName", + "options": "error" + }, + "properties": [ + { + "id": "custom.width", + "value": 900 + } + ] + }, + { + "matcher": { + "id": "byName", + "options": "job" + }, + "properties": [ + { + "id": "links", + "value": [ + { + "title": "Filter dashboard to this job", + "url": "/d/gpkg-merger-logs?${__url_time_range}&var-loki=${loki}&${namespace:queryparam}&var-job=${__data.fields.job}&var-task=", + "targetBlank": false + } + ] + } + ] + }, + { + "matcher": { + "id": "byName", + "options": "task" + }, + "properties": [ + { + "id": "links", + "value": [ + { + "title": "Filter dashboard to this job and task", + "url": "/d/gpkg-merger-logs?${__url_time_range}&var-loki=${loki}&${namespace:queryparam}&var-job=${__data.fields.job}&var-task=${__data.fields.task}", + "targetBlank": false + } + ] + } + ] + } + ] + }, + "options": { + "showHeader": true, + "cellHeight": "sm", + "footer": { + "show": false + } + }, + "targets": [ + { + "refId": "A", + "datasource": { + "type": "loki", + "uid": "${loki}" + }, + "editorMode": "code", + "expr": "{k8s_namespace_name=~\"$namespace\", k8s_deployment_name=~\"gpkg-merger.*\"} | json | __error__=\"\" | jobId=~\"$job\" | taskId=~\"$task\" | level=~\"Error|Critical\"", + "queryType": "range" + } + ], + "transformations": [ + { + "id": "extractFields", + "options": { + "source": "labels", + "replace": false, + "keepTime": false + } + }, + { + "id": "organize", + "options": { + "excludeByName": { + "service": true, + "version": true, + "thread": true, + "service_name": true, + "k8s_container_name": true, + "k8s_deployment_name": true, + "k8s_namespace_name": true, + "k8s_pod_name": true, + "time": true, + "Line": true, + "labels": true, + "id": true, + "tsNs": true, + "category": true, + "level": true, + "added": true, + "merged": true, + "replaced": true, + "skipped": true, + "total": true, + "durationSeconds": true + }, + "indexByName": { + "Time": 0, + "jobId": 1, + "taskId": 2, + "message": 3 + }, + "renameByName": { + "jobId": "job", + "taskId": "task", + "message": "error" + } + } + } + ] + }, + { + "id": 14, + "type": "table", + "title": "Log timeline (parsed)", + "datasource": { + "type": "loki", + "uid": "${loki}" + }, + "gridPos": { + "h": 10, + "w": 24, + "x": 0, + "y": 42 + }, + "fieldConfig": { + "defaults": { + "custom": { + "filterable": true, + "align": "auto" + } + }, + "overrides": [] + }, + "options": { + "showHeader": true, + "cellHeight": "sm", + "footer": { + "show": false + } + }, + "targets": [ + { + "refId": "A", + "datasource": { + "type": "loki", + "uid": "${loki}" + }, + "editorMode": "code", + "expr": "{k8s_namespace_name=~\"$namespace\", k8s_deployment_name=~\"gpkg-merger.*\"} | json | __error__=\"\" | jobId=~\"$job\" | taskId=~\"$task\"", + "queryType": "range" + } + ], + "transformations": [ + { + "id": "extractFields", + "options": { + "source": "labels", + "replace": false, + "keepTime": false + } + }, + { + "id": "organize", + "options": { + "excludeByName": { + "service": true, + "version": true, + "thread": true, + "service_name": true, + "k8s_container_name": true, + "k8s_deployment_name": true, + "k8s_namespace_name": true, + "k8s_pod_name": true, + "time": true, + "Line": true, + "labels": true, + "id": true, + "tsNs": true, + "jobId": true, + "taskId": true, + "added": true, + "merged": true, + "replaced": true, + "skipped": true, + "total": true, + "durationSeconds": true + }, + "indexByName": { + "Time": 0, + "level": 1, + "category": 2, + "message": 3 + }, + "renameByName": {} + } + } + ] + } + ], + "annotations": { + "list": [] + } +} \ No newline at end of file diff --git a/MergerLogic/Extensions/ServiceCollectionExtensions.cs b/MergerLogic/Extensions/ServiceCollectionExtensions.cs index e3eeda97..087e0527 100644 --- a/MergerLogic/Extensions/ServiceCollectionExtensions.cs +++ b/MergerLogic/Extensions/ServiceCollectionExtensions.cs @@ -152,6 +152,8 @@ public static IServiceCollection RegisterOpenTelemetry(this IServiceCollection c options.AddProcessor( new SimpleLogRecordExportProcessor(new OpenTelemetryFormattedConsoleExporter(new ConsoleExporterOptions()))); //lgtm [cs/local-not-disposed] options.SetResourceBuilder(resourceBuilder); + // required for the exporter to emit BeginScope correlation fields (jobId/taskId) + options.IncludeScopes = true; }); }); #endregion Logger diff --git a/MergerLogic/ImageProcessing/ITileMerger.cs b/MergerLogic/ImageProcessing/ITileMerger.cs index b1bd4610..749072c3 100644 --- a/MergerLogic/ImageProcessing/ITileMerger.cs +++ b/MergerLogic/ImageProcessing/ITileMerger.cs @@ -6,5 +6,7 @@ namespace MergerLogic.ImageProcessing public interface ITileMerger { Tile? MergeTiles(List tiles, Coord targetCoords, TileFormatStrategy strategy, bool uploadOnly = false); + + Tile? MergeTiles(List tiles, Coord targetCoords, TileFormatStrategy strategy, out MergeStats stats, bool uploadOnly = false); } } diff --git a/MergerLogic/ImageProcessing/MergeStats.cs b/MergerLogic/ImageProcessing/MergeStats.cs new file mode 100644 index 00000000..4fee673d --- /dev/null +++ b/MergerLogic/ImageProcessing/MergeStats.cs @@ -0,0 +1,18 @@ +namespace MergerLogic.ImageProcessing +{ + /// + /// Describes which inputs contributed to a merged tile, used to classify the + /// write as added / merged / replaced. TargetUsed is false in upload-only mode. + /// + public readonly struct MergeStats + { + public bool TargetUsed { get; } + public bool AnySourceUsed { get; } + + public MergeStats(bool targetUsed, bool anySourceUsed) + { + this.TargetUsed = targetUsed; + this.AnySourceUsed = anySourceUsed; + } + } +} diff --git a/MergerLogic/ImageProcessing/TileMerger.cs b/MergerLogic/ImageProcessing/TileMerger.cs index 4a9a5a24..0f0ca14e 100644 --- a/MergerLogic/ImageProcessing/TileMerger.cs +++ b/MergerLogic/ImageProcessing/TileMerger.cs @@ -20,6 +20,15 @@ public TileMerger(ITileScaler tileScaler, ILogger logger) public Tile? MergeTiles(List tiles, Coord targetCoords, TileFormatStrategy strategy, bool uploadOnly = false) { + return this.MergeTiles(tiles, targetCoords, strategy, out _, uploadOnly); + } + + public Tile? MergeTiles(List tiles, Coord targetCoords, TileFormatStrategy strategy, + out MergeStats stats, bool uploadOnly = false) + { + bool targetUsed = false; + bool anySourceUsed = false; + if(uploadOnly) { this._logger.LogDebug($"[{MethodBase.GetCurrentMethod()?.Name}] Configured to upload only mode"); // Ignore target if in upload only mode @@ -30,11 +39,13 @@ public TileMerger(ITileScaler tileScaler, ILogger logger) this._logger.LogDebug($"[{MethodBase.GetCurrentMethod()?.Name}] Only one source was found, using raw image"); Tile? rawTile = tiles[0](); rawTile?.ConvertToFormat(strategy.ApplyStrategy(rawTile.Format)); + stats = new MergeStats(false, rawTile != null); return rawTile; } } - var images = this.GetImageList(tiles, targetCoords, uploadOnly); + bool hasTarget = !uploadOnly && tiles.Count > 0; + var images = this.GetImageList(tiles, targetCoords, uploadOnly, hasTarget, out targetUsed, out anySourceUsed); IMagickImage image; switch (images.Count) @@ -42,6 +53,7 @@ public TileMerger(ITileScaler tileScaler, ILogger logger) case 0: // There are no images this._logger.LogDebug($"[{MethodBase.GetCurrentMethod()?.Name}] No images where found return null"); + stats = new MergeStats(targetUsed, anySourceUsed); return null; case 1: ImageFormatter.RemoveImageDateAttributes(images[0]); @@ -73,14 +85,18 @@ public TileMerger(ITileScaler tileScaler, ILogger logger) Tile tile = new Tile(targetCoords, image); image.Dispose(); tile.ConvertToFormat(strategy.ApplyStrategy(tile.Format)); + stats = new MergeStats(targetUsed, anySourceUsed); return tile; } - private List GetImageList(List tiles, Coord targetCoords, bool uploadOnly) + private List GetImageList(List tiles, Coord targetCoords, bool uploadOnly, + bool hasTarget, out bool targetUsed, out bool anySourceUsed) { var images = new List(); int i = tiles.Count - 1; Tile? tile = null; + targetUsed = false; + anySourceUsed = false; bool hasAlpha = false; try @@ -100,7 +116,20 @@ private List GetImageList(List tiles, Coo continue; } + int before = images.Count; this.AddTileToImageList(targetCoords, tile, images, out hasAlpha); + if (images.Count > before) + { + if (hasTarget && i == 0) + { + targetUsed = true; + } + else + { + anySourceUsed = true; + } + } + if (!hasAlpha) { return images; diff --git a/MergerLogic/Monitoring/OpenTelemetryFormattedConsoleExporter.cs b/MergerLogic/Monitoring/OpenTelemetryFormattedConsoleExporter.cs index 45c56f97..908cc4b2 100644 --- a/MergerLogic/Monitoring/OpenTelemetryFormattedConsoleExporter.cs +++ b/MergerLogic/Monitoring/OpenTelemetryFormattedConsoleExporter.cs @@ -1,6 +1,7 @@ -using OpenTelemetry; +using OpenTelemetry; using OpenTelemetry.Exporter; using OpenTelemetry.Logs; +using System.Text.Json; namespace MergerLogic.Monitoring { @@ -29,9 +30,44 @@ private string MCTextFormat(LogRecord record) var resource = this.ParseResource(); var serviceName = this.GetResourceAttribute(resource, SERVICE_NAME_ATTRIBUTE, "unknown_service"); var serviceVersion = this.GetResourceAttribute(resource, SERVICE_VERSION_ATTRIBUTE, "unknown_version"); - var exception = record.Exception != null ? $" [{record.Exception}]" : string.Empty; - return $"[{this.FormatTime(record.Timestamp)}] [{record.LogLevel}] [{serviceName}] [{serviceVersion}] [{record.CategoryName}] [{Environment.CurrentManagedThreadId}] {record.State}{exception}"; + var entry = new Dictionary + { + ["time"] = this.FormatTime(record.Timestamp), + ["level"] = record.LogLevel.ToString(), + ["service"] = serviceName, + ["version"] = serviceVersion, + ["category"] = record.CategoryName, + ["thread"] = Environment.CurrentManagedThreadId, + ["message"] = record.State?.ToString(), + }; + + if (record.Exception != null) + { + entry["exception"] = record.Exception.ToString(); + } + + this.AddScopes(record, entry); + + return JsonSerializer.Serialize(entry); + } + + // Flatten ILogger.BeginScope key/value pairs to top-level fields (e.g. jobId/taskId) so they + // are queryable in Loki. Empty unless IncludeScopes is enabled (see DI setup). + private void AddScopes(LogRecord record, Dictionary entry) + { + record.ForEachScope((scope, state) => + { + foreach (var pair in scope) + { + if (pair.Key == "{OriginalFormat}") + { + continue; + } + + state[pair.Key] = pair.Value; + } + }, entry); } private string FormatTime(DateTime time) diff --git a/MergerLogicUnitTests/ImageProcessing/TileMergerTest.cs b/MergerLogicUnitTests/ImageProcessing/TileMergerTest.cs index d8cde850..1be3cc96 100644 --- a/MergerLogicUnitTests/ImageProcessing/TileMergerTest.cs +++ b/MergerLogicUnitTests/ImageProcessing/TileMergerTest.cs @@ -339,6 +339,65 @@ public void MergeTiles(Tile[] tiles, Coord targetCoord, TileFormatStrategy strat CollectionAssert.AreEqual(expectedTileBytes, result.GetImageBytes()); } + [TestMethod] + [TestCategory("MergeTiles")] + public void MergeTilesStatsBlendedTargetAndSource() + { + var targetCoord = new Coord(15, 0, 0); + // target (index 0) is transparent, source (last) is transparent -> both enter the stack + var tiles = new[] + { + new Tile(targetCoord, File.ReadAllBytes("2.png")), + new Tile(targetCoord, File.ReadAllBytes("1.png")) + }; + var tileBuilders = tiles.Select(tile => () => tile).ToList(); + + var result = this._testTileMerger.MergeTiles(tileBuilders, targetCoord, new TileFormatStrategy(TileFormat.Png), out var stats); + + Assert.IsNotNull(result); + Assert.IsTrue(stats.TargetUsed); + Assert.IsTrue(stats.AnySourceUsed); + } + + [TestMethod] + [TestCategory("MergeTiles")] + public void MergeTilesStatsOpaqueSourceOverTarget() + { + var targetCoord = new Coord(15, 0, 0); + // opaque source (last) short-circuits before the target (index 0) is reached + var tiles = new[] + { + new Tile(targetCoord, File.ReadAllBytes("1.png")), + new Tile(targetCoord, File.ReadAllBytes("3.jpeg")) + }; + var tileBuilders = tiles.Select(tile => () => tile).ToList(); + + var result = this._testTileMerger.MergeTiles(tileBuilders, targetCoord, new TileFormatStrategy(TileFormat.Jpeg), out var stats); + + Assert.IsNotNull(result); + Assert.IsFalse(stats.TargetUsed); + Assert.IsTrue(stats.AnySourceUsed); + } + + [TestMethod] + [TestCategory("MergeTiles")] + public void MergeTilesStatsUploadOnly() + { + var targetCoord = new Coord(15, 0, 0); + var tiles = new[] + { + new Tile(targetCoord, File.ReadAllBytes("2.png")), + new Tile(targetCoord, File.ReadAllBytes("1.png")) + }; + var tileBuilders = tiles.Select(tile => () => tile).ToList(); + + var result = this._testTileMerger.MergeTiles(tileBuilders, targetCoord, new TileFormatStrategy(TileFormat.Jpeg), out var stats, uploadOnly: true); + + Assert.IsNotNull(result); + Assert.IsFalse(stats.TargetUsed); + Assert.IsTrue(stats.AnySourceUsed); + } + #endregion } } diff --git a/MergerLogicUnitTests/Monitoring/OpenTelemetryFormattedConsoleExporterTest.cs b/MergerLogicUnitTests/Monitoring/OpenTelemetryFormattedConsoleExporterTest.cs new file mode 100644 index 00000000..9336b2cb --- /dev/null +++ b/MergerLogicUnitTests/Monitoring/OpenTelemetryFormattedConsoleExporterTest.cs @@ -0,0 +1,92 @@ +using MergerLogic.Monitoring; +using Microsoft.Extensions.Logging; +using Microsoft.VisualStudio.TestTools.UnitTesting; +using OpenTelemetry; +using OpenTelemetry.Exporter; +using OpenTelemetry.Logs; +using System; +using System.Collections.Generic; +using System.IO; +using System.Text.Json; + +namespace MergerLogicUnitTests.Monitoring +{ + [TestClass] + [TestCategory("unit")] + [TestCategory("monitoring")] + public class OpenTelemetryFormattedConsoleExporterTest + { + private TextWriter _originalOut = null!; + + [TestInitialize] + public void BeforeEach() + { + this._originalOut = Console.Out; + } + + [TestCleanup] + public void AfterEach() + { + Console.SetOut(this._originalOut); + } + + private static ILoggerFactory BuildFactory() + { + return LoggerFactory.Create(builder => + { + builder.ClearProviders(); + builder.AddOpenTelemetry(options => + { + options.IncludeScopes = true; + options.AddProcessor(new SimpleLogRecordExportProcessor( + new OpenTelemetryFormattedConsoleExporter(new ConsoleExporterOptions()))); + }); + }); + } + + private static JsonElement CaptureSingleLine(Action log) + { + var writer = new StringWriter(); + Console.SetOut(writer); + + using (var factory = BuildFactory()) + { + log(factory.CreateLogger("TestCategory")); + } + + string line = writer.ToString().Trim(); + return JsonDocument.Parse(line).RootElement; + } + + [TestMethod] + public void WhenLoggingInsideAScope_ShouldEmitScopeAsTopLevelJsonFields() + { + JsonElement entry = CaptureSingleLine(logger => + { + using (logger.BeginScope(new Dictionary + { + ["jobId"] = "job-123", + ["taskId"] = "task-456" + })) + { + logger.LogInformation("processing tiles"); + } + }); + + Assert.AreEqual("processing tiles", entry.GetProperty("message").GetString()); + Assert.AreEqual("job-123", entry.GetProperty("jobId").GetString()); + Assert.AreEqual("task-456", entry.GetProperty("taskId").GetString()); + Assert.AreEqual("Information", entry.GetProperty("level").GetString()); + } + + [TestMethod] + public void WhenLoggingWithoutAScope_ShouldNotEmitScopeFields() + { + JsonElement entry = CaptureSingleLine(logger => logger.LogInformation("no scope here")); + + Assert.AreEqual("no scope here", entry.GetProperty("message").GetString()); + Assert.IsFalse(entry.TryGetProperty("jobId", out _), "unexpected jobId field"); + Assert.IsFalse(entry.TryGetProperty("taskId", out _), "unexpected taskId field"); + } + } +} diff --git a/MergerService/Models/Jobs/JobParamersAdditiomalParams.cs b/MergerService/Models/Jobs/JobParamersAdditiomalParams.cs index 9ef7b236..f2126701 100644 --- a/MergerService/Models/Jobs/JobParamersAdditiomalParams.cs +++ b/MergerService/Models/Jobs/JobParamersAdditiomalParams.cs @@ -7,13 +7,15 @@ namespace MergerService.Models.Jobs public class AdditionalParams { [JsonInclude] public string? JobTrackerServiceURL { get; } + [JsonInclude] public string? ReportOutputPath { get; } [System.Text.Json.Serialization.JsonIgnore] private JsonSerializerSettings _jsonSerializerSettings; - public AdditionalParams(string jobTrackerServiceURL) + public AdditionalParams(string jobTrackerServiceURL, string? reportOutputPath = null) { this.JobTrackerServiceURL = jobTrackerServiceURL; + this.ReportOutputPath = reportOutputPath; this._jsonSerializerSettings = new JsonSerializerSettings(); this._jsonSerializerSettings.Converters.Add(new StringEnumConverter()); diff --git a/MergerService/Models/Reports/MergeReport.cs b/MergerService/Models/Reports/MergeReport.cs new file mode 100644 index 00000000..ad892f61 --- /dev/null +++ b/MergerService/Models/Reports/MergeReport.cs @@ -0,0 +1,119 @@ +using MergerLogic.DataTypes; +using MergerLogic.ImageProcessing; +using Newtonsoft.Json; +using Newtonsoft.Json.Serialization; + +namespace MergerService.Models.Reports +{ + public class MergeReport + { + public int Version => 1; + public string JobId { get; } + public string TaskId { get; } + public string TaskType { get; } + public string TargetFormat { get; } + public bool IsNewTarget { get; } + + public DateTime StartTime { get; private set; } + public DateTime EndTime { get; private set; } + public double DurationSeconds { get; private set; } + + public int Added { get; private set; } + public int Merged { get; private set; } + public int Replaced { get; private set; } + public int Skipped { get; private set; } + public int Total => this.Added + this.Merged + this.Replaced + this.Skipped; + + public double AddedPercentage { get; private set; } + public double MergedPercentage { get; private set; } + public double ReplacedPercentage { get; private set; } + public double SkippedPercentage { get; private set; } + + [JsonProperty("addedTiles")] + public List AddedTiles { get; } = new List(); + + public MergeReport(string jobId, string taskId, string taskType, string targetFormat, bool isNewTarget) + { + this.JobId = jobId; + this.TaskId = taskId; + this.TaskType = taskType; + this.TargetFormat = targetFormat; + this.IsNewTarget = isNewTarget; + } + + // Classifies a single merged tile into added / merged / replaced / skipped. + // added = tile didn't exist in target before the merge; merged = existed and the + // target was blended with source data; replaced = existed and an opaque source + // covered it; skipped = no tile produced or no source data contributed. + public void RecordOutcome(Coord coord, bool existedBefore, bool tileProduced, MergeStats stats) + { + if (!tileProduced || !stats.AnySourceUsed) + { + this.RecordSkipped(); + return; + } + + if (!existedBefore) + { + this.RecordAdded(coord); + return; + } + + if (stats.TargetUsed) + { + this.RecordMerged(); + return; + } + + this.RecordReplaced(); + } + + private void RecordAdded(Coord coord) + { + this.Added++; + // store a copy so a later in-place coord mutation can't corrupt the list + this.AddedTiles.Add(new Coord(coord.Z, coord.X, coord.Y)); + } + + private void RecordMerged() => this.Merged++; + private void RecordReplaced() => this.Replaced++; + private void RecordSkipped() => this.Skipped++; + + public void Finalize(DateTime startTime, DateTime endTime) + { + this.StartTime = startTime; + this.EndTime = endTime; + this.DurationSeconds = (endTime - startTime).TotalSeconds; + + int total = this.Total; + if (total > 0) + { + this.AddedPercentage = 100.0 * this.Added / total; + this.MergedPercentage = 100.0 * this.Merged / total; + this.ReplacedPercentage = 100.0 * this.Replaced / total; + this.SkippedPercentage = 100.0 * this.Skipped / total; + } + } + + public string ToJson() => JsonConvert.SerializeObject(this, Formatting.None, + new JsonSerializerSettings { ContractResolver = new CamelCasePropertyNamesContractResolver() }); + + // Summary for the structured log line: counts + percentages, WITHOUT the added-tile list. + public string ToLogString() + { + var summary = new + { + version = this.Version, + jobId = this.JobId, + taskId = this.TaskId, + taskType = this.TaskType, + targetFormat = this.TargetFormat, + isNewTarget = this.IsNewTarget, + durationSeconds = this.DurationSeconds, + counts = new { added = this.Added, merged = this.Merged, replaced = this.Replaced, skipped = this.Skipped, total = this.Total }, + percentages = new { added = this.AddedPercentage, merged = this.MergedPercentage, replaced = this.ReplacedPercentage, skipped = this.SkippedPercentage } + }; + return JsonConvert.SerializeObject(summary, Formatting.None); + } + } +} diff --git a/MergerService/Program.cs b/MergerService/Program.cs index e327200c..7379ecb4 100644 --- a/MergerService/Program.cs +++ b/MergerService/Program.cs @@ -29,6 +29,7 @@ builder.Services.AddSingleton(); builder.Services.AddSingleton(); builder.Services.AddSingleton(); +builder.Services.AddSingleton(); builder.Services.AddSingleton(); var app = builder.Build(); diff --git a/MergerService/Runners/ITaskExecutor.cs b/MergerService/Runners/ITaskExecutor.cs index 1a2cdd8a..503995c8 100644 --- a/MergerService/Runners/ITaskExecutor.cs +++ b/MergerService/Runners/ITaskExecutor.cs @@ -5,6 +5,6 @@ namespace MergerService.Runners { public interface ITaskExecutor { - void ExecuteTask(MergeTask task, ITaskUtils taskUtils, string? managerCallbackUrl); + void ExecuteTask(MergeTask task, ITaskUtils taskUtils, string? managerCallbackUrl, string? reportOutputPath); } } diff --git a/MergerService/Runners/TaskExecutor.cs b/MergerService/Runners/TaskExecutor.cs index b63b1a3c..e71f854a 100644 --- a/MergerService/Runners/TaskExecutor.cs +++ b/MergerService/Runners/TaskExecutor.cs @@ -4,6 +4,7 @@ using MergerLogic.Monitoring.Metrics; using MergerLogic.Utils; using MergerService.Controllers; +using MergerService.Models.Reports; using MergerService.Models.Tasks; using MergerService.Utils; using System.Diagnostics; @@ -21,6 +22,7 @@ public class TaskExecutor : ITaskExecutor private readonly ActivitySource _activitySource; private readonly IFileSystem _fileSystem; private readonly IMetricsProvider _metricsProvider; + private readonly IReportWriter _reportWriter; private readonly string _inputPath; private readonly string _gpkgPath; private readonly bool _limitBatchSize; @@ -32,7 +34,7 @@ public class TaskExecutor : ITaskExecutor public TaskExecutor(IDataFactory dataFactory, ITileMerger tileMerger, ITimeUtils timeUtils, IConfigurationManager configurationManager, ILogger logger, ActivitySource activitySource, - IFileSystem fileSystem, IMetricsProvider metricsProvider) + IFileSystem fileSystem, IMetricsProvider metricsProvider, IReportWriter reportWriter) { this._dataFactory = dataFactory; this._tileMerger = tileMerger; @@ -41,6 +43,7 @@ public TaskExecutor(IDataFactory dataFactory, ITileMerger tileMerger, ITimeUtils this._activitySource = activitySource; this._fileSystem = fileSystem; this._metricsProvider = metricsProvider; + this._reportWriter = reportWriter; this._inputPath = configurationManager.GetConfiguration("GENERAL", "inputPath"); this._gpkgPath = configurationManager.GetConfiguration("GENERAL", "gpkgPath"); this._filePath = configurationManager.GetConfiguration("GENERAL", "filePath"); @@ -61,7 +64,7 @@ public TaskExecutor(IDataFactory dataFactory, ITileMerger tileMerger, ITimeUtils } } - public void ExecuteTask(MergeTask task, ITaskUtils taskUtils, string? managerCallbackUrl) + public void ExecuteTask(MergeTask task, ITaskUtils taskUtils, string? managerCallbackUrl, string? reportOutputPath) { string methodName = MethodBase.GetCurrentMethod().Name; this._logger.LogDebug($"[{methodName}] start {task.ToString()}"); @@ -83,6 +86,9 @@ public void ExecuteTask(MergeTask task, ITaskUtils taskUtils, string? managerCal } MergeMetadata metadata = task.Parameters; + DateTime reportStart = DateTime.UtcNow; + MergeReport report = new MergeReport(task.JobId, task.Id, task.Type, + metadata.TargetFormat.ToString(), metadata.IsNewTarget); Stopwatch mergeRunTimeStopwatch = new Stopwatch(); TimeSpan ts; @@ -162,11 +168,17 @@ public void ExecuteTask(MergeTask task, ITaskUtils taskUtils, string? managerCal // TODO: upscale = false - this is a temporary fix till we decide how sources should be upscaled correspondingTileBuilders.Add(() => source.GetCorrespondingTile(coord, false)); } + // TileExists mutates coord.Y in place; pass a copy so the coord the + // merge builders capture stays in its original grid/origin space. + bool existedBefore = !metadata.IsNewTarget && target.TileExists(new Coord(coord.Z, coord.X, coord.Y)); + var tileMergeStopwatch = Stopwatch.StartNew(); - Tile? tile = this._tileMerger.MergeTiles(correspondingTileBuilders, coord, strategy, metadata.IsNewTarget); + Tile? tile = this._tileMerger.MergeTiles(correspondingTileBuilders, coord, strategy, out MergeStats stats, metadata.IsNewTarget); tileMergeStopwatch.Stop(); this._metricsProvider.MergeTimePerTileHistogram(tileMergeStopwatch.Elapsed.TotalSeconds, metadata.TargetFormat); + report.RecordOutcome(coord, existedBefore, tile != null, stats); + if (tile != null) { tiles.Add(tile); @@ -253,6 +265,33 @@ public void ExecuteTask(MergeTask task, ITaskUtils taskUtils, string? managerCal } target.Wrapup(); } + + report.Finalize(reportStart, DateTime.UtcNow); + // emit the report counts as structured scope fields (top-level in the JSON log line) + // so per-job/task statistics can be aggregated straight from the logs + using (this._logger.BeginScope(new Dictionary + { + ["added"] = report.Added, + ["merged"] = report.Merged, + ["replaced"] = report.Replaced, + ["skipped"] = report.Skipped, + ["total"] = report.Total, + ["durationSeconds"] = report.DurationSeconds, + })) + { + this._logger.LogInformation($"[{methodName}] Merge report: {report.ToLogString()}"); + } + try + { + this._reportWriter.WriteReport(report, reportOutputPath); + } + catch (Exception e) + { + // Best-effort (proposed default). Whether this should fail the task is an open + // question raised on the implementation PR. + this._logger.LogError(e, $"[{methodName}] Failed to write merge report artifact: {e.Message}"); + } + this._logger.LogDebug($"[{methodName}] end"); } diff --git a/MergerService/Runners/TaskRunner.cs b/MergerService/Runners/TaskRunner.cs index 00884162..ef86376f 100644 --- a/MergerService/Runners/TaskRunner.cs +++ b/MergerService/Runners/TaskRunner.cs @@ -1,5 +1,6 @@ using MergerLogic.Clients; using MergerLogic.Monitoring.Metrics; +using MergerService.Models.Jobs; using MergerService.Models.Tasks; using MergerService.Utils; using System.Diagnostics; @@ -86,73 +87,79 @@ public bool RunTask(MergeTask? task) return false; } - this._logger.LogInformation($"[{methodName}] Run Task: jobId {task.JobId}, taskId {task.Id}"); - string? managerCallbackUrl = this._jobUtils.GetJob(task.JobId)?.Parameters.AdditionalParams?.JobTrackerServiceURL; - string log = managerCallbackUrl == null ? "managerCallbackUrl not provided as job parameter" : $"managerCallback url: {managerCallbackUrl}"; - this._logger.LogDebug($"[{methodName}]{log}"); - - // check if needs to fail task that was released by task liberator and reached max attempts - if (task.Attempts >= this._maxTaskRetriesAttempts) + // tag every log line from this task with jobId/taskId for per-task correlation + using (this._logger.BeginScope(new Dictionary { ["jobId"] = task.JobId, ["taskId"] = task.Id })) { + this._logger.LogInformation($"[{methodName}] Run Task: jobId {task.JobId}, taskId {task.Id}"); + MergeJob? job = this._jobUtils.GetJob(task.JobId); + string? managerCallbackUrl = job?.Parameters.AdditionalParams?.JobTrackerServiceURL; + string? reportOutputPath = job?.Parameters.AdditionalParams?.ReportOutputPath; + string log = managerCallbackUrl == null ? "managerCallbackUrl not provided as job parameter" : $"managerCallback url: {managerCallbackUrl}"; + this._logger.LogDebug($"[{methodName}]{log}"); + + // check if needs to fail task that was released by task liberator and reached max attempts + if (task.Attempts >= this._maxTaskRetriesAttempts) + { + try + { + string reason = string.IsNullOrEmpty(task.Reason) ? $"Max attempts reached, current attempt is {task.Attempts}" : $"{task.Reason} and Max attempts reached with {task.Attempts} attempts"; + this._logger.LogWarning($"[{methodName}] reject job because attemts count reached, jobId {task.JobId}, taskId {task.Id}, {reason}"); + this._taskUtils.UpdateReject(task.JobId, task.Id, task.Attempts, reason, task.Resettable, managerCallbackUrl); + } + catch (Exception innerError) + { + this._logger.LogError(innerError, $"[{methodName}] Error in MergerService while updating reject status for job {task.JobId}, task {task.Id} due to max attemps reached with {task.Attempts}, update task failure: {innerError.Message}"); + } + + return false; + } + + var totalTaskStopwatch = Stopwatch.StartNew(); + bool taskSucceed = false; + try { - string reason = string.IsNullOrEmpty(task.Reason) ? $"Max attempts reached, current attempt is {task.Attempts}" : $"{task.Reason} and Max attempts reached with {task.Attempts} attempts"; - this._logger.LogWarning($"[{methodName}] reject job because attemts count reached, jobId {task.JobId}, taskId {task.Id}, {reason}"); - this._taskUtils.UpdateReject(task.JobId, task.Id, task.Attempts, reason, task.Resettable, managerCallbackUrl); + this._heartbeatClient.Start(task.Id); + this._taskExecutor.ExecuteTask(task, this._taskUtils, managerCallbackUrl, reportOutputPath); + taskSucceed = true; } - catch (Exception innerError) + catch (Exception e) { - this._logger.LogError(innerError, $"[{methodName}] Error in MergerService while updating reject status for job {task.JobId}, task {task.Id} due to max attemps reached with {task.Attempts}, update task failure: {innerError.Message}"); + this._logger.LogError(e, $"[{methodName}] Error in MergerService while running task {task.Id}, error: {e.Message}"); + + try + { + this._taskUtils.UpdateReject(task.JobId, task.Id, task.Attempts, e.Message, true, managerCallbackUrl); + } + catch (Exception innerError) + { + this._logger.LogError(e, $"[{methodName}] Error in MergerService while updating reject status, RunTask catch block - update task failure: {innerError.Message}"); + } + } + finally + { + totalTaskStopwatch.Stop(); + this._metricsProvider.TaskExecutionTimeHistogram(totalTaskStopwatch.Elapsed.TotalSeconds, task.Type); + this._heartbeatClient.Stop(); } - return false; - } - - var totalTaskStopwatch = Stopwatch.StartNew(); - bool taskSucceed = false; - - try - { - this._heartbeatClient.Start(task.Id); - this._taskExecutor.ExecuteTask(task, this._taskUtils, managerCallbackUrl); - taskSucceed = true; - } - catch (Exception e) - { - this._logger.LogError(e, $"[{methodName}] Error in MergerService while running task {task.Id}, error: {e.Message}"); + if (!taskSucceed) + { + return false; + } try { - this._taskUtils.UpdateReject(task.JobId, task.Id, task.Attempts, e.Message, true, managerCallbackUrl); + this._taskUtils.UpdateCompletion(task.JobId, task.Id, managerCallbackUrl); + this._logger.LogInformation($"[{methodName}] Completed task: jobId: {task.JobId}, taskId: {task.Id}"); } - catch (Exception innerError) + catch (Exception e) { - this._logger.LogError(e, $"[{methodName}] Error in MergerService while updating reject status, RunTask catch block - update task failure: {innerError.Message}"); + this._logger.LogError(e, $"[{methodName}] Error in MergerService start - update task completion: {e.Message}"); } - } - finally - { - totalTaskStopwatch.Stop(); - this._metricsProvider.TaskExecutionTimeHistogram(totalTaskStopwatch.Elapsed.TotalSeconds, task.Type); - this._heartbeatClient.Stop(); - } - if (!taskSucceed) - { - return false; + return true; } - - try - { - this._taskUtils.UpdateCompletion(task.JobId, task.Id, managerCallbackUrl); - this._logger.LogInformation($"[{methodName}] Completed task: jobId: {task.JobId}, taskId: {task.Id}"); - } - catch (Exception e) - { - this._logger.LogError(e, $"[{methodName}] Error in MergerService start - update task completion: {e.Message}"); - } - - return true; } } } diff --git a/MergerService/Utils/IReportWriter.cs b/MergerService/Utils/IReportWriter.cs new file mode 100644 index 00000000..31f0f003 --- /dev/null +++ b/MergerService/Utils/IReportWriter.cs @@ -0,0 +1,11 @@ +using MergerService.Models.Reports; + +namespace MergerService.Utils +{ + public interface IReportWriter + { + // Writes the report JSON artifact to the configured sink under outputPath. + // No-op when outputPath is null/empty. Throws on write failure. + void WriteReport(MergeReport report, string? outputPath); + } +} diff --git a/MergerService/Utils/ReportWriter.cs b/MergerService/Utils/ReportWriter.cs new file mode 100644 index 00000000..7de2af97 --- /dev/null +++ b/MergerService/Utils/ReportWriter.cs @@ -0,0 +1,75 @@ +using Amazon.S3; +using Amazon.S3.Model; +using MergerLogic.Utils; +using MergerService.Models.Reports; +using Microsoft.Extensions.DependencyInjection; +using Microsoft.Extensions.Logging; +using System.IO.Abstractions; +using System.Reflection; + +namespace MergerService.Utils +{ + public class ReportWriter : IReportWriter + { + private readonly IConfigurationManager _configuration; + private readonly IFileSystem _fileSystem; + private readonly IServiceProvider _serviceProvider; + private readonly ILogger _logger; + + public ReportWriter(IConfigurationManager configuration, IFileSystem fileSystem, IServiceProvider serviceProvider, + ILogger logger) + { + this._configuration = configuration; + this._fileSystem = fileSystem; + this._serviceProvider = serviceProvider; + this._logger = logger; + } + + public void WriteReport(MergeReport report, string? outputPath) + { + string methodName = MethodBase.GetCurrentMethod().Name; + if (string.IsNullOrEmpty(outputPath)) + { + this._logger.LogDebug($"[{methodName}] No ReportOutputPath configured, skipping report artifact"); + return; + } + + string fileName = $"merge-report-{report.JobId}-{report.TaskId}.json"; + string json = report.ToJson(); + string sink = this._configuration.GetConfiguration("REPORT", "sink"); + + if (string.Equals(sink, "S3", System.StringComparison.OrdinalIgnoreCase)) + { + this.WriteToS3(outputPath, fileName, json); + } + else + { + this.WriteToFs(outputPath, fileName, json); + } + + this._logger.LogInformation($"[{methodName}] Wrote merge report to {sink}:{outputPath}/{fileName}"); + } + + private void WriteToFs(string outputPath, string fileName, string json) + { + this._fileSystem.Directory.CreateDirectory(outputPath); + string fullPath = this._fileSystem.Path.Combine(outputPath, fileName); + this._fileSystem.File.WriteAllText(fullPath, json); + } + + private void WriteToS3(string outputPath, string fileName, string json) + { + string bucket = this._configuration.GetConfiguration("S3", "bucket"); + string key = $"{outputPath.TrimEnd('/')}/{fileName}"; + var request = new PutObjectRequest + { + BucketName = bucket, + Key = key, + ContentBody = json, + ContentType = "application/json" + }; + var s3 = this._serviceProvider.GetRequiredService(); + s3.PutObjectAsync(request).Wait(); + } + } +} diff --git a/MergerService/appsettings.json b/MergerService/appsettings.json index 1d96c4c1..21929561 100644 --- a/MergerService/appsettings.json +++ b/MergerService/appsettings.json @@ -52,6 +52,9 @@ "port": 9500, "measurementBuckets": "[0.001, 0.005, 0.01, 0.025, 0.05, 0.1, 0.25, 0.5, 1, 2.5, 5, 10, 15, 50,250, 500]" }, + "REPORT": { + "sink": "FS" + }, "AllowedHosts": "*", "HTTP": { "retries": 3 diff --git a/MergerServiceUnitTests/MergerServiceUnitTests.csproj b/MergerServiceUnitTests/MergerServiceUnitTests.csproj index 4213e99a..77426672 100644 --- a/MergerServiceUnitTests/MergerServiceUnitTests.csproj +++ b/MergerServiceUnitTests/MergerServiceUnitTests.csproj @@ -12,6 +12,7 @@ + all runtime; build; native; contentfiles; analyzers; buildtransitive diff --git a/MergerServiceUnitTests/Models/MergeReportTest.cs b/MergerServiceUnitTests/Models/MergeReportTest.cs new file mode 100644 index 00000000..30636835 --- /dev/null +++ b/MergerServiceUnitTests/Models/MergeReportTest.cs @@ -0,0 +1,119 @@ +using MergerLogic.DataTypes; +using MergerLogic.ImageProcessing; +using MergerService.Models.Reports; +using Microsoft.VisualStudio.TestTools.UnitTesting; +using Newtonsoft.Json.Linq; +using System; + +namespace MergerServiceUnitTests.Models +{ + [TestClass] + [TestCategory("unit")] + public class MergeReportTest + { + private static readonly Coord AnyCoord = new Coord(10, 1, 2); + + [TestMethod] + public void RecordOutcome_TileNotProduced_CountsAsSkipped() + { + var report = new MergeReport("job1", "task1", "MERGE", "PNG", false); + report.RecordOutcome(AnyCoord, existedBefore: true, tileProduced: false, new MergeStats(true, true)); + Assert.AreEqual(1, report.Skipped); + Assert.AreEqual(0, report.Added + report.Merged + report.Replaced); + } + + [TestMethod] + public void RecordOutcome_NoSourceData_CountsAsSkipped() + { + var report = new MergeReport("job1", "task1", "MERGE", "PNG", false); + report.RecordOutcome(AnyCoord, existedBefore: false, tileProduced: true, new MergeStats(true, false)); + Assert.AreEqual(1, report.Skipped); + Assert.AreEqual(0, report.Added); + } + + [TestMethod] + public void RecordOutcome_NotExistedBefore_CountsAsAdded() + { + var report = new MergeReport("job1", "task1", "MERGE", "PNG", false); + report.RecordOutcome(new Coord(10, 5, 6), existedBefore: false, tileProduced: true, new MergeStats(false, true)); + Assert.AreEqual(1, report.Added); + Assert.AreEqual(1, report.AddedTiles.Count); + Assert.AreEqual(new Coord(10, 5, 6), report.AddedTiles[0]); + } + + [TestMethod] + public void RecordOutcome_ExistedAndTargetBlended_CountsAsMerged() + { + var report = new MergeReport("job1", "task1", "MERGE", "PNG", false); + report.RecordOutcome(AnyCoord, existedBefore: true, tileProduced: true, new MergeStats(true, true)); + Assert.AreEqual(1, report.Merged); + Assert.AreEqual(0, report.Added); + } + + [TestMethod] + public void RecordOutcome_ExistedAndOpaqueSource_CountsAsReplaced() + { + var report = new MergeReport("job1", "task1", "MERGE", "PNG", false); + report.RecordOutcome(AnyCoord, existedBefore: true, tileProduced: true, new MergeStats(false, true)); + Assert.AreEqual(1, report.Replaced); + Assert.AreEqual(0, report.Merged); + } + + [TestMethod] + public void RecordOutcome_StoresCopyOfAddedCoord() + { + var report = new MergeReport("job1", "task1", "MERGE", "PNG", false); + var coord = new Coord(10, 7, 8); + report.RecordOutcome(coord, existedBefore: false, tileProduced: true, new MergeStats(false, true)); + + // mutating the caller's coord must not affect the stored one + coord.Y = 999; + Assert.AreEqual(8, report.AddedTiles[0].Y); + } + + [TestMethod] + public void Counts_And_Percentages_Are_Computed() + { + var report = new MergeReport("job1", "task1", "MERGE", "PNG", false); + report.RecordOutcome(new Coord(10, 1, 2), existedBefore: false, tileProduced: true, new MergeStats(false, true)); // added + report.RecordOutcome(new Coord(10, 1, 3), existedBefore: false, tileProduced: true, new MergeStats(false, true)); // added + report.RecordOutcome(AnyCoord, existedBefore: true, tileProduced: true, new MergeStats(true, true)); // merged + report.RecordOutcome(AnyCoord, existedBefore: true, tileProduced: true, new MergeStats(false, true)); // replaced + report.RecordOutcome(AnyCoord, existedBefore: false, tileProduced: false, new MergeStats(false, false)); // skipped + + report.Finalize(new DateTime(2026, 9, 15, 0, 0, 0, DateTimeKind.Utc), + new DateTime(2026, 9, 15, 0, 1, 0, DateTimeKind.Utc)); + + Assert.AreEqual(2, report.Added); + Assert.AreEqual(1, report.Merged); + Assert.AreEqual(1, report.Replaced); + Assert.AreEqual(1, report.Skipped); + Assert.AreEqual(5, report.Total); + Assert.AreEqual(60, report.DurationSeconds); + Assert.AreEqual(40.0, report.AddedPercentage, 0.01); + } + + [TestMethod] + public void Json_Includes_AddedTiles_LogString_Excludes_Them() + { + var report = new MergeReport("job1", "task1", "MERGE", "PNG", false); + report.RecordOutcome(new Coord(10, 1, 2), existedBefore: false, tileProduced: true, new MergeStats(false, true)); + report.Finalize(DateTime.UnixEpoch, DateTime.UnixEpoch); + + JObject json = JObject.Parse(report.ToJson()); + Assert.AreEqual(1, ((JArray)json["addedTiles"]).Count); + + Assert.IsFalse(report.ToLogString().Contains("addedTiles")); + Assert.IsTrue(report.ToLogString().Contains("\"added\"")); + } + + [TestMethod] + public void Percentages_Are_Zero_When_No_Tiles() + { + var report = new MergeReport("job1", "task1", "MERGE", "PNG", true); + report.Finalize(DateTime.UnixEpoch, DateTime.UnixEpoch); + Assert.AreEqual(0, report.Total); + Assert.AreEqual(0.0, report.AddedPercentage, 0.01); + } + } +} diff --git a/MergerServiceUnitTests/Runners/TaskExecutorTest.cs b/MergerServiceUnitTests/Runners/TaskExecutorTest.cs index 5522ac56..8e27d804 100644 --- a/MergerServiceUnitTests/Runners/TaskExecutorTest.cs +++ b/MergerServiceUnitTests/Runners/TaskExecutorTest.cs @@ -4,6 +4,7 @@ using MergerLogic.Monitoring.Metrics; using MergerLogic.Utils; using MergerService.Controllers; +using MergerService.Models.Reports; using MergerService.Models.Tasks; using MergerService.Runners; using MergerService.Utils; @@ -37,6 +38,8 @@ public class TaskExecutorTest private Mock _taskUtilsMock; private Mock _tileScalerMock; private Mock> _tileMergerLoggerMock; + private Mock _reportWriterMock; + private Mock _tileMergerMock; private ActivitySource _testActivitySource; private ITileMerger _testTileMerger; @@ -67,6 +70,8 @@ public void BeforeEach() this._taskUtilsMock = this._mockRepository.Create(); this._tileScalerMock = this._mockRepository.Create(); this._tileMergerLoggerMock = this._mockRepository.Create>(); + this._reportWriterMock = this._mockRepository.Create(); + this._tileMergerMock = this._mockRepository.Create(); this._testActivitySource = new ActivitySource("test"); this._testTileMerger = new TileMerger(_tileScalerMock.Object, _tileMergerLoggerMock.Object); @@ -89,13 +94,13 @@ public void WhenGivenSourcesWithOneTile_ShouldWriteAllTilesToTarget(int numberOf var testTaskExecutor = new TaskExecutor(_dataFactoryMock.Object, _testTileMerger, _timeUtilsMock.Object, _configurationManagerMock.Object, _taskExecutorLoggerMock.Object, _testActivitySource, _testFileSystem, - _metricsProviderMock.Object); + _metricsProviderMock.Object, _reportWriterMock.Object); targetDataMock.Setup(targetData => targetData.UpdateTiles(It.IsAny>())).Callback>( resultWrittenTiles.AddRange ); - testTaskExecutor.ExecuteTask(testTask, _taskUtilsMock.Object, null); + testTaskExecutor.ExecuteTask(testTask, _taskUtilsMock.Object, null, null); targetDataMock.Verify(targetData => targetData.UpdateTiles(It.Is>( tiles => tiles.All( @@ -174,13 +179,13 @@ public void WhenConfiguringBatchLimits_ShouldWriteTilesEachTimeAfterReachingBatc var resultWrittenTiles = new List(); var testTaskExecutor = new TaskExecutor(_dataFactoryMock.Object, _testTileMerger, _timeUtilsMock.Object, _configurationManagerMock.Object, _taskExecutorLoggerMock.Object, _testActivitySource, _testFileSystem, - _metricsProviderMock.Object); + _metricsProviderMock.Object, _reportWriterMock.Object); targetDataMock.Setup(targetData => targetData.UpdateTiles(It.IsAny>())).Callback>( resultWrittenTiles.AddRange ); - testTaskExecutor.ExecuteTask(testTask, _taskUtilsMock.Object, null); + testTaskExecutor.ExecuteTask(testTask, _taskUtilsMock.Object, null, null); targetDataMock.Verify(targetData => targetData.UpdateTiles(It.IsAny>()), Times.Exactly(amountOfFlushes)); targetDataMock.Verify(targetData => targetData.Wrapup(), Times.Once); @@ -190,6 +195,65 @@ public void WhenConfiguringBatchLimits_ShouldWriteTilesEachTimeAfterReachingBatc )); } + // Integration/wiring test: verifies ExecuteTask feeds RecordOutcome the right inputs + // (existedBefore, tile produced, MergeStats) and writes the resulting report. + // The exhaustive added/merged/replaced/skipped classification permutations live in MergeReportTest. + [TestMethod] + [TestCategory("unit")] + [TestCategory("runners")] + public void ExecuteTask_ClassifiesTileAndWritesReport() + { + this._configurationManagerMock.Setup(configManager => configManager.GetConfiguration("GENERAL", "batchSize", "batchMaxSize")).Returns(1); + this._configurationManagerMock.Setup(configManager => configManager.GetConfiguration("GENERAL", "batchSize", "limitBatchSize")).Returns(true); + this._configurationManagerMock.Setup(configManager => configManager.GetConfiguration("GENERAL", "batchMaxBytes")).Returns(1); + + byte[] tileBytes = File.ReadAllBytes("tile.jpeg"); + Source testTarget = new Source("target", "target_type", new Extent(), GridOrigin.UPPER_LEFT, Grid.TwoXOne); + Source testSource = new Source("source", "source_type"); + Mock targetDataMock = this._mockRepository.Create(); + Mock sourceDataMock = this._mockRepository.Create(); + + this._dataFactoryMock.Setup(dataFactory => dataFactory.CreateDataSource( + testTarget.Type, testTarget.Path, It.IsAny(), + testTarget.Grid, testTarget.Origin, testTarget.Extent, It.IsAny()) + ).Returns(targetDataMock.Object); + this._dataFactoryMock.Setup(dataFactory => dataFactory.CreateDataSource( + testSource.Type, testSource.Path, It.IsAny(), + testSource.Grid, testSource.Origin, testSource.Extent, It.IsAny()) + ).Returns(sourceDataMock.Object); + + // target already has this tile → existedBefore == true; merger reports target blended → merged + targetDataMock.Setup(targetData => targetData.TileExists(It.IsAny())).Returns(true); + + TileBounds tileBounds = new TileBounds(1, 1, 1, 1, 1); + var testTask = new MergeTask("id", "type", "description", + new MergeMetadata(TileFormat.Jpeg, false, new TileBounds[] { tileBounds }, new Source[] { testTarget, testSource }), + Status.PENDING, null, "reason", 0, "jobId", true, new DateTime(), new DateTime()); + + MergeStats outStats = new MergeStats(targetUsed: true, anySourceUsed: true); + this._tileMergerMock.Setup(m => m.MergeTiles( + It.IsAny>(), It.IsAny(), + It.IsAny(), out outStats, It.IsAny()) + ).Returns(new Tile(new Coord(1, 1, 1), tileBytes)); + + MergeReport captured = null; + this._reportWriterMock.Setup(w => w.WriteReport(It.IsAny(), "reports")) + .Callback((r, p) => captured = r); + + var testTaskExecutor = new TaskExecutor(_dataFactoryMock.Object, _tileMergerMock.Object, _timeUtilsMock.Object, + _configurationManagerMock.Object, _taskExecutorLoggerMock.Object, _testActivitySource, _testFileSystem, + _metricsProviderMock.Object, _reportWriterMock.Object); + + testTaskExecutor.ExecuteTask(testTask, _taskUtilsMock.Object, null, "reports"); + + this._reportWriterMock.Verify(w => w.WriteReport(It.IsAny(), "reports"), Times.Once); + Assert.IsNotNull(captured); + Assert.AreEqual(1, captured.Merged); + Assert.AreEqual(0, captured.Added); + Assert.AreEqual(0, captured.Replaced); + Assert.AreEqual(0, captured.Skipped); + } + private Tuple, Tile[]> SetupTestTask(int amountOfSources, bool isTargetNew) { byte[] tileBytes = File.ReadAllBytes("tile.jpeg"); diff --git a/MergerServiceUnitTests/Runners/TaskRunnerTest.cs b/MergerServiceUnitTests/Runners/TaskRunnerTest.cs index c518ee45..e0754cee 100644 --- a/MergerServiceUnitTests/Runners/TaskRunnerTest.cs +++ b/MergerServiceUnitTests/Runners/TaskRunnerTest.cs @@ -3,6 +3,7 @@ using MergerLogic.ImageProcessing; using MergerLogic.Monitoring.Metrics; using MergerService.Controllers; +using MergerService.Models.Jobs; using MergerService.Models.Tasks; using MergerService.Runners; using MergerService.Utils; @@ -72,8 +73,14 @@ public void WhenTaskExecutedSuccessfully_ShouldUpdateTaskCompletion() new MergeMetadata(TileFormat.Jpeg, true, new TileBounds[0], new Source[0]), Status.PENDING, 0, "reason", 0, "testJobId", true, new DateTime(), new DateTime()); + var testJob = new MergeJob("testJobId", "resourceId", "version", "type", "resolution", "description", + new JobMergeMetadata(null!, new string[0], "", "", new AdditionalParams("http://tracker", "reports")), + new DateTime(), new DateTime(), Status.PENDING, 0, "reason", false, 0, "internalId", "producerName", + "productName", "productType", 0, 0, 0, 0, 0, 0, 0, "additionalIdentifiers", "domain", new MergeTask[0]); + this._taskUtilsMock.Setup(taskUtils => taskUtils.GetTask(It.IsAny(), It.IsAny())).Returns(testTask); - this._taskExecutorMock.Setup(taskExecutor => taskExecutor.ExecuteTask(testTask, _taskUtilsMock.Object, It.IsAny())); + this._jobUtilsMock.Setup(jobUtils => jobUtils.GetJob(testTask.JobId)).Returns(testJob); + this._taskExecutorMock.Setup(taskExecutor => taskExecutor.ExecuteTask(testTask, _taskUtilsMock.Object, It.IsAny(), It.IsAny())); var testTaskRunner = new TaskRunner(_taskExecutorMock.Object, _jobUtilsMock.Object, _loggerMock.Object, _taskUtilsMock.Object, _heartbeatClientMock.Object, _metricsProviderMock.Object, @@ -83,9 +90,35 @@ public void WhenTaskExecutedSuccessfully_ShouldUpdateTaskCompletion() testTaskRunner.RunTask(testResultTask); Assert.AreEqual(testTask, testResultTask); + this._taskExecutorMock.Verify(e => e.ExecuteTask(testTask, _taskUtilsMock.Object, It.IsAny(), "reports"), Times.Once); _taskUtilsMock.Verify(taskUtils => taskUtils.UpdateCompletion(testTask.JobId, testTask.Id, It.IsAny()), Times.Once); } + [TestMethod] + public void WhenRunningTask_ShouldScopeLogsWithJobAndTaskId() + { + var testTask = new MergeTask("testTaskId", "type", "description", + new MergeMetadata(TileFormat.Jpeg, true, new TileBounds[0], new Source[0]), + Status.PENDING, 0, "reason", 0, "testJobId", true, new DateTime(), new DateTime()); + + Dictionary? capturedScope = null; + this._taskUtilsMock.Setup(taskUtils => taskUtils.GetTask(It.IsAny(), It.IsAny())).Returns(testTask); + this._taskExecutorMock.Setup(taskExecutor => taskExecutor.ExecuteTask(testTask, _taskUtilsMock.Object, It.IsAny(), It.IsAny())); + this._loggerMock.Setup(logger => logger.BeginScope(It.IsAny>())) + .Returns(Mock.Of()) + .Callback>(scope => capturedScope = scope); + + var testTaskRunner = new TaskRunner(_taskExecutorMock.Object, _jobUtilsMock.Object, _loggerMock.Object, + _taskUtilsMock.Object, _heartbeatClientMock.Object, _metricsProviderMock.Object, + _configurationManagerMock.Object); + + testTaskRunner.RunTask(testTaskRunner.FetchTask(new KeyValuePair("testJobType", "testTaskType"))); + + Assert.IsNotNull(capturedScope); + Assert.AreEqual(testTask.JobId, capturedScope["jobId"]); + Assert.AreEqual(testTask.Id, capturedScope["taskId"]); + } + [TestMethod] public void WhenTaskExecutionFailed_ShouldUpdateTaskFailed() { @@ -95,7 +128,7 @@ public void WhenTaskExecutionFailed_ShouldUpdateTaskFailed() Status.PENDING, 0, "reason", 0, "testJobId", true, new DateTime(), new DateTime()); this._taskUtilsMock.Setup(taskUtils => taskUtils.GetTask(It.IsAny(), It.IsAny())).Returns(testTask); - this._taskExecutorMock.Setup(taskExecutor => taskExecutor.ExecuteTask(testTask, _taskUtilsMock.Object, It.IsAny())).Throws(new Exception(testFailureMessage)); + this._taskExecutorMock.Setup(taskExecutor => taskExecutor.ExecuteTask(testTask, _taskUtilsMock.Object, It.IsAny(), It.IsAny())).Throws(new Exception(testFailureMessage)); var testTaskRunner = new TaskRunner(_taskExecutorMock.Object, _jobUtilsMock.Object, _loggerMock.Object, _taskUtilsMock.Object, _heartbeatClientMock.Object, _metricsProviderMock.Object, diff --git a/MergerServiceUnitTests/Utils/ReportWriterTest.cs b/MergerServiceUnitTests/Utils/ReportWriterTest.cs new file mode 100644 index 00000000..64ef4e55 --- /dev/null +++ b/MergerServiceUnitTests/Utils/ReportWriterTest.cs @@ -0,0 +1,89 @@ +using Amazon.S3; +using Amazon.S3.Model; +using MergerLogic.Utils; +using MergerService.Models.Reports; +using MergerService.Utils; +using Microsoft.Extensions.DependencyInjection; +using Microsoft.Extensions.Logging; +using Microsoft.VisualStudio.TestTools.UnitTesting; +using Moq; +using System; +using System.IO.Abstractions; +using System.IO.Abstractions.TestingHelpers; +using System.Linq; +using System.Threading; +using System.Threading.Tasks; + +namespace MergerServiceUnitTests.Utils +{ + [TestClass] + [TestCategory("unit")] + public class ReportWriterTest + { + private Mock _config; + private Mock> _logger; + private Mock _s3; + private Mock _serviceProvider; + + [TestInitialize] + public void BeforeEach() + { + this._config = new Mock(MockBehavior.Loose); + this._logger = new Mock>(MockBehavior.Loose); + this._s3 = new Mock(MockBehavior.Loose); + this._serviceProvider = new Mock(MockBehavior.Loose); + this._serviceProvider.Setup(sp => sp.GetService(typeof(IAmazonS3))).Returns(this._s3.Object); + } + + private MergeReport BuildReport() + { + var r = new MergeReport("job1", "task1", "MERGE", "PNG", false); + r.Finalize(System.DateTime.UnixEpoch, System.DateTime.UnixEpoch); + return r; + } + + [TestMethod] + public void FsSink_WritesJsonFile() + { + this._config.Setup(c => c.GetConfiguration("REPORT", "sink")).Returns("FS"); + var fs = new MockFileSystem(); + var writer = new ReportWriter(this._config.Object, fs, this._serviceProvider.Object, this._logger.Object); + + writer.WriteReport(this.BuildReport(), "/reports"); + + string expected = fs.Path.Combine("/reports", "merge-report-job1-task1.json"); + Assert.IsTrue(fs.FileExists(expected)); + StringAssert.Contains(fs.File.ReadAllText(expected), "\"jobId\":\"job1\""); + } + + [TestMethod] + public void S3Sink_PutsObject() + { + this._config.Setup(c => c.GetConfiguration("REPORT", "sink")).Returns("S3"); + this._config.Setup(c => c.GetConfiguration("S3", "bucket")).Returns("tiles"); + this._s3.Setup(s => s.PutObjectAsync(It.IsAny(), It.IsAny())) + .ReturnsAsync(new PutObjectResponse()); + var fs = new MockFileSystem(); + var writer = new ReportWriter(this._config.Object, fs, this._serviceProvider.Object, this._logger.Object); + + writer.WriteReport(this.BuildReport(), "reports"); + + this._s3.Verify(s => s.PutObjectAsync( + It.Is(r => r.BucketName == "tiles" && r.Key == "reports/merge-report-job1-task1.json"), + It.IsAny()), Times.Once); + } + + [TestMethod] + public void EmptyPath_IsNoOp() + { + var fs = new MockFileSystem(); + var writer = new ReportWriter(this._config.Object, fs, this._serviceProvider.Object, this._logger.Object); + + writer.WriteReport(this.BuildReport(), null); + writer.WriteReport(this.BuildReport(), ""); + + Assert.AreEqual(0, fs.AllFiles.Count()); + this._s3.Verify(s => s.PutObjectAsync(It.IsAny(), It.IsAny()), Times.Never); + } + } +}