Skip to content

Commit 1d6be74

Browse files
authored
Expose pipeline outcome classes (#2173)
* fix: expose pipeline outcome classes * fix: align outcome class assignments
1 parent 7501f47 commit 1d6be74

5 files changed

Lines changed: 202 additions & 4 deletions

File tree

inc/Abilities/Engine/ExecuteStepAbility.php

Lines changed: 6 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -246,6 +246,7 @@ public function execute( array $input ): array {
246246
'step_success' => $step_success,
247247
'packet_count' => count( $dataPackets ),
248248
'status' => $recorded_status,
249+
'reason' => $result['reason'] ?? ( $result['error'] ?? null ),
249250
'error' => $result['error'] ?? null,
250251
)
251252
);
@@ -595,6 +596,7 @@ private function routeAfterExecution(
595596
'success' => true,
596597
'step_success' => false,
597598
'outcome' => 'failed',
599+
'reason' => $transition_failure_reason,
598600
'error' => $transition_failure_reason,
599601
);
600602
}
@@ -742,21 +744,24 @@ private function routeAfterExecution(
742744
'step_type' => $step_type,
743745
)
744746
);
747+
$empty_packet_reason = $this->getFailureReasonFromPackets( $dataPackets, 'empty_data_packet_returned' );
748+
745749
do_action(
746750
'datamachine_fail_job',
747751
$job_id,
748752
'step_execution_failure',
749753
array(
750754
'flow_step_id' => $flow_step_id,
751755
'class' => $step_class,
752-
'reason' => $this->getFailureReasonFromPackets( $dataPackets, 'empty_data_packet_returned' ),
756+
'reason' => $empty_packet_reason,
753757
)
754758
);
755759

756760
return array(
757761
'success' => true,
758762
'step_success' => false,
759763
'outcome' => 'failed',
764+
'reason' => $empty_packet_reason,
760765
);
761766
}
762767

inc/Cli/Commands/JobsCommand.php

Lines changed: 5 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -1777,6 +1777,7 @@ public function metrics( array $args, array $assoc_args ): void {
17771777
$counts = $metrics['counts'] ?? array();
17781778
$children = $metrics['child_jobs'] ?? array();
17791779
$timestamps = $metrics['timestamps'] ?? array();
1780+
$classes = is_array( $metrics['outcome_classes'] ?? null ) ? $metrics['outcome_classes'] : array();
17801781

17811782
WP_CLI::log( sprintf( 'Job ID: %d', $metrics['job_id'] ?? 0 ) );
17821783
WP_CLI::log( sprintf( 'Status: %s', $metrics['status'] ?? '' ) );
@@ -1796,6 +1797,10 @@ public function metrics( array $args, array $assoc_args ): void {
17961797
}
17971798
WP_CLI::log( '' );
17981799

1800+
WP_CLI::log( 'Outcome Classes:' );
1801+
WP_CLI::log( empty( $classes ) ? ' -' : ' ' . implode( ', ', array_map( 'strval', $classes ) ) );
1802+
WP_CLI::log( '' );
1803+
17991804
WP_CLI::log( 'Child Jobs:' );
18001805
foreach ( $children as $key => $value ) {
18011806
WP_CLI::log( sprintf( ' %s: %d', $key, (int) $value ) );

inc/Core/RunMetrics.php

Lines changed: 124 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -23,7 +23,14 @@ class RunMetrics {
2323
'failed',
2424
'fetch_packets',
2525
'no_content',
26+
'true_empty_query',
27+
'provider_error',
28+
'hydration_failed',
29+
'hydration_partial',
30+
'ai_empty_packet',
31+
'missing_handler_packet',
2632
'source_rejected',
33+
'item_deferred',
2734
'retried',
2835
'scheduled',
2936
'staged_actions',
@@ -66,6 +73,12 @@ public static function recordStepResult( int $job_id, string $flow_step_id, arra
6673
if ( 'source_rejected' === ( $clean_result['result'] ?? '' ) || ! empty( $clean_result['source_rejection_reason'] ) ) {
6774
$metrics['counts']['source_rejected'] = max( 1, (int) $metrics['counts']['source_rejected'] );
6875
}
76+
foreach ( self::classesFromStepResult( $clean_result, (string) ( $clean_result['status'] ?? '' ) ) as $class ) {
77+
if ( ! isset( $metrics['counts'][ $class ] ) ) {
78+
$metrics['counts'][ $class ] = 0;
79+
}
80+
$metrics['counts'][ $class ] = max( 1, (int) $metrics['counts'][ $class ] );
81+
}
6982
$metrics['last_activity_at'] = self::now();
7083
$engine[ self::KEY ] = $metrics;
7184

@@ -161,6 +174,8 @@ public static function fromJob( array $job ): array {
161174
$counts[ $key ] = max( (int) ( $counts[ $key ] ?? 0 ), (int) $value );
162175
}
163176

177+
$outcome_classes = self::outcomeClasses( $job, $engine, $counts );
178+
164179
return array(
165180
'job_id' => (int) ( $job['job_id'] ?? 0 ),
166181
'source' => $job['source'] ?? null,
@@ -170,6 +185,7 @@ public static function fromJob( array $job ): array {
170185
'parent_job_id' => isset( $job['parent_job_id'] ) ? (int) $job['parent_job_id'] : 0,
171186
'status' => $status,
172187
'counts' => $counts,
188+
'outcome_classes' => $outcome_classes,
173189
'child_jobs' => self::childTotals( (int) ( $job['job_id'] ?? 0 ) ),
174190
'timestamps' => array(
175191
'created_at' => $job['created_at'] ?? null,
@@ -178,7 +194,7 @@ public static function fromJob( array $job ): array {
178194
'completed_at' => $ended_at,
179195
),
180196
'duration_seconds' => self::durationSeconds( $started_at, $duration_end ),
181-
'outcome' => self::outcomeDetails( $job, $engine, $counts ),
197+
'outcome' => self::outcomeDetails( $job, $engine, $counts, $outcome_classes ),
182198
'step_results' => self::stepResults( $engine ),
183199
'context' => $metrics['context'],
184200
'token_usage' => self::tokenUsage( $engine ),
@@ -270,22 +286,36 @@ private static function inferCountsFromEngine( array $engine, string $status ):
270286
if ( 'source_rejected' === ( $step_result['result'] ?? '' ) || ! empty( $step_result['source_rejection_reason'] ) ) {
271287
$counts['source_rejected'] = max( 1, (int) ( $counts['source_rejected'] ?? 0 ) );
272288
}
289+
foreach ( self::classesFromStepResult( $step_result, $status ) as $class ) {
290+
$counts[ $class ] = max( 1, (int) ( $counts[ $class ] ?? 0 ) );
291+
}
292+
}
293+
294+
foreach ( self::classesFromStatus( $status ) as $class ) {
295+
$counts[ $class ] = max( 1, (int) ( $counts[ $class ] ?? 0 ) );
273296
}
274297

275298
return $counts;
276299
}
277300

278-
private static function outcomeDetails( array $job, array $engine, array $counts ): array {
301+
private static function outcomeDetails( array $job, array $engine, array $counts, array $outcome_classes ): array {
279302
$status = (string) ( $job['status'] ?? '' );
280303
$job_status = JobStatus::fromString( $status );
281304
$source_rejection = is_array( $engine['source_rejection'] ?? null ) ? $engine['source_rejection'] : array();
282305
$is_source_rejected = ! empty( $counts['source_rejected'] ) || ! empty( $source_rejection ) || 'source-rejected' === $job_status->getReason();
306+
$class_counts = array();
307+
foreach ( $outcome_classes as $class ) {
308+
$class_counts[ $class ] = (int) ( $counts[ $class ] ?? 0 );
309+
}
283310

284311
return array_filter(
285312
array(
286313
'status' => $status,
287314
'base_status' => $job_status->getBaseStatus(),
288315
'status_reason' => $job_status->getReason(),
316+
'primary_class' => $outcome_classes[0] ?? null,
317+
'classes' => $outcome_classes,
318+
'class_counts' => $class_counts,
289319
'fetch_packet_count' => (int) ( $counts['fetch_packets'] ?? 0 ),
290320
'no_content' => ! empty( $counts['no_content'] ) || JobStatus::COMPLETED_NO_ITEMS === $job_status->getBaseStatus(),
291321
'source_rejected' => $is_source_rejected,
@@ -298,6 +328,98 @@ private static function outcomeDetails( array $job, array $engine, array $counts
298328
);
299329
}
300330

331+
private static function outcomeClasses( array $job, array $engine, array $counts ): array {
332+
$status = (string) ( $job['status'] ?? '' );
333+
$classes = self::classesFromStatus( $status );
334+
335+
foreach ( self::stepResults( $engine ) as $step_result ) {
336+
$classes = array_merge( $classes, self::classesFromStepResult( $step_result, $status ) );
337+
}
338+
339+
foreach ( array( 'true_empty_query', 'provider_error', 'hydration_failed', 'hydration_partial', 'ai_empty_packet', 'missing_handler_packet', 'source_rejected', 'item_deferred' ) as $class ) {
340+
if ( ! empty( $counts[ $class ] ) ) {
341+
$classes[] = $class;
342+
}
343+
}
344+
345+
return array_values( array_unique( array_filter( $classes ) ) );
346+
}
347+
348+
private static function classesFromStepResult( array $step_result, string $status = '' ): array {
349+
$result = (string) ( $step_result['result'] ?? '' );
350+
$reason = self::outcomeReasonFrom( $step_result, $status );
351+
$classes = self::classesFromReason( $reason );
352+
353+
if ( in_array( $result, array( 'no_content', 'completed_no_items' ), true ) ) {
354+
$classes[] = 'true_empty_query';
355+
}
356+
if ( 'source_rejected' === $result || ! empty( $step_result['source_rejection_reason'] ) ) {
357+
$classes[] = 'source_rejected';
358+
}
359+
if ( 'fetch' === ( $step_result['step_type'] ?? '' ) && 'failed' === $result ) {
360+
$classes[] = 'provider_error';
361+
}
362+
363+
return array_values( array_unique( array_filter( $classes ) ) );
364+
}
365+
366+
private static function classesFromStatus( string $status ): array {
367+
$job_status = JobStatus::fromString( $status );
368+
$classes = self::classesFromReason( (string) $job_status->getReason() );
369+
370+
if ( JobStatus::COMPLETED_NO_ITEMS === $job_status->getBaseStatus() ) {
371+
$classes[] = 'true_empty_query';
372+
}
373+
374+
return array_values( array_unique( array_filter( $classes ) ) );
375+
}
376+
377+
private static function classesFromReason( string $reason ): array {
378+
$reason = strtolower( str_replace( '-', '_', trim( $reason ) ) );
379+
if ( '' === $reason ) {
380+
return array();
381+
}
382+
383+
$classes = array();
384+
if ( in_array( $reason, array( 'mcp_fetch_failed', 'auth_ref_resolution_failed', 'ai_provider_missing' ), true ) || str_contains( $reason, 'provider' ) ) {
385+
$classes[] = 'provider_error';
386+
}
387+
if ( 'missing_source_content' === $reason || str_contains( $reason, 'hydration_failed' ) ) {
388+
$classes[] = 'hydration_failed';
389+
}
390+
if ( str_contains( $reason, 'hydration_partial' ) ) {
391+
$classes[] = 'hydration_partial';
392+
}
393+
if ( 'empty_data_packet_returned' === $reason ) {
394+
$classes[] = 'ai_empty_packet';
395+
}
396+
if ( 'handler_requiring_step_missing_handler_packets' === $reason ) {
397+
$classes[] = 'missing_handler_packet';
398+
}
399+
if ( 'source_rejected' === $reason ) {
400+
$classes[] = 'source_rejected';
401+
}
402+
if ( 'item_deferred' === $reason ) {
403+
$classes[] = 'item_deferred';
404+
}
405+
406+
return $classes;
407+
}
408+
409+
private static function outcomeReasonFrom( array $step_result, string $status ): string {
410+
foreach ( array( 'reason', 'error', 'status' ) as $key ) {
411+
$value = $step_result[ $key ] ?? null;
412+
if ( is_scalar( $value ) && '' !== (string) $value ) {
413+
if ( 'status' === $key ) {
414+
return (string) JobStatus::fromString( (string) $value )->getReason();
415+
}
416+
return (string) $value;
417+
}
418+
}
419+
420+
return (string) JobStatus::fromString( $status )->getReason();
421+
}
422+
301423
private static function stepResults( array $engine ): array {
302424
$step_results = is_array( $engine[ self::STEP_RESULTS_KEY ] ?? null ) ? $engine[ self::STEP_RESULTS_KEY ] : array();
303425
return array_values( array_filter( $step_results, 'is_array' ) );

inc/Core/Steps/Fetch/Tools/FetchItemDispositionTool.php

Lines changed: 15 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -245,6 +245,21 @@ private function deferItem( array $parameters, array $tool_def ): array {
245245
)
246246
);
247247

248+
if ( $flow_step_id && class_exists( RunMetrics::class ) ) {
249+
RunMetrics::recordStepResult(
250+
$job_id,
251+
(string) $flow_step_id,
252+
array(
253+
'step_type' => 'fetch',
254+
'result' => 'item_deferred',
255+
'packet_count' => 0,
256+
'reason' => 'item-deferred',
257+
'source_type' => $source_type,
258+
'item_identifier' => $item_identifier,
259+
)
260+
);
261+
}
262+
248263
do_action(
249264
'datamachine_log',
250265
'info',

tests/job-outcome-metrics-smoke.php

Lines changed: 52 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -88,6 +88,8 @@
8888
);
8989
$assert( 'base status is completed_no_items', 'completed_no_items' === $metrics['outcome']['base_status'] );
9090
$assert( 'no_content boolean is true', true === $metrics['outcome']['no_content'] );
91+
$assert( 'true empty query class is exposed', in_array( 'true_empty_query', $metrics['outcome_classes'], true ) );
92+
$assert( 'true empty query count is exposed', 1 === $metrics['counts']['true_empty_query'] );
9193

9294
echo "\n[3] failed status exposes failure base status\n";
9395
$metrics = RunMetrics::fromJob( $job( 'failed - api-timeout', array() ) );
@@ -113,14 +115,63 @@
113115
);
114116
$assert( 'source_rejected boolean is true', true === $metrics['outcome']['source_rejected'] );
115117
$assert( 'source rejection reason is exposed', 'not relevant' === $metrics['outcome']['source_rejection_reason'] );
118+
$assert( 'source rejected outcome class is exposed', array( 'source_rejected' ) === $metrics['outcome_classes'] );
116119

117-
echo "\n[5] CLI/source integration markers exist\n";
120+
echo "\n[5] failure reasons expose distinct generic outcome classes\n";
121+
$metrics = RunMetrics::fromJob( $job( 'failed - mcp_fetch_failed', array() ) );
122+
$assert( 'provider failure class is exposed from status reason', array( 'provider_error' ) === $metrics['outcome_classes'] );
123+
$assert( 'provider failure count is exposed', 1 === $metrics['counts']['provider_error'] );
124+
125+
$metrics = RunMetrics::fromJob( $job( 'failed - missing_source_content', array() ) );
126+
$assert( 'hydration failure class is exposed from status reason', array( 'hydration_failed' ) === $metrics['outcome_classes'] );
127+
128+
$metrics = RunMetrics::fromJob(
129+
$job(
130+
'failed - step_execution_failure',
131+
array(
132+
'step_results' => array(
133+
'ai_1' => array(
134+
'flow_step_id' => 'ai_1',
135+
'step_type' => 'ai',
136+
'result' => 'failed',
137+
'reason' => 'empty_data_packet_returned',
138+
'packet_count' => 0,
139+
),
140+
),
141+
)
142+
)
143+
);
144+
$assert( 'AI empty packet class is exposed from step result', array( 'ai_empty_packet' ) === $metrics['outcome_classes'] );
145+
146+
$metrics = RunMetrics::fromJob(
147+
$job(
148+
'failed - step_execution_failure',
149+
array(
150+
'step_results' => array(
151+
'ai_1' => array(
152+
'flow_step_id' => 'ai_1',
153+
'step_type' => 'ai',
154+
'result' => 'failed',
155+
'reason' => 'handler_requiring_step_missing_handler_packets',
156+
'packet_count' => 2,
157+
),
158+
),
159+
)
160+
)
161+
);
162+
$assert( 'missing handler packet class is exposed from step result', array( 'missing_handler_packet' ) === $metrics['outcome_classes'] );
163+
164+
$metrics = RunMetrics::fromJob( $job( 'failed - item-deferred', array() ) );
165+
$assert( 'item deferred class is exposed from status reason', array( 'item_deferred' ) === $metrics['outcome_classes'] );
166+
167+
echo "\n[6] CLI/source integration markers exist\n";
118168
$jobs_command = file_get_contents( __DIR__ . '/../inc/Cli/Commands/JobsCommand.php' ) ?: '';
119169
$fetch_step = file_get_contents( __DIR__ . '/../inc/Core/Steps/Fetch/FetchStep.php' ) ?: '';
120170
$disposition = file_get_contents( __DIR__ . '/../inc/Core/Steps/Fetch/Tools/FetchItemDispositionTool.php' ) ?: '';
121171
$assert( 'jobs list supports pipeline filter', str_contains( $jobs_command, "assoc_args['pipeline']" ) );
122172
$assert( 'jobs list supports handler filter', str_contains( $jobs_command, "assoc_args['handler']" ) );
123173
$assert( 'jobs list JSON includes outcome', str_contains( $jobs_command, "item['outcome']" ) );
174+
$assert( 'jobs metrics table prints outcome classes', str_contains( $jobs_command, 'Outcome Classes:' ) );
124175
$assert( 'fetch step records packet count', str_contains( $fetch_step, "'packet_count' => count( \$packets )" ) );
125176
$assert( 'source rejection persists structured reason', str_contains( $disposition, "'source_rejection'" ) );
126177

0 commit comments

Comments
 (0)