Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
8 changes: 7 additions & 1 deletion CLAUDE.md
Original file line number Diff line number Diff line change
Expand Up @@ -162,7 +162,7 @@ background process on port 8080.
| `workflow start --sync` | Start and wait for completion | None | `--workflow`, `--input`, `--file`, `--wait-until` | `conductor workflow start --workflow my_workflow --sync` |
| `workflow status <id>` | Get execution status | workflow ID | | `conductor workflow status abc-123` |
| `workflow get-execution <id>` | Get full execution details | workflow ID | `--complete` | `conductor workflow get-execution abc-123` |
| `workflow search` | Search executions | None | `--workflow`, `--status`, `--count`, `--start-time-after`, `--start-time-before`, `--json` | `conductor workflow search --workflow my_workflow --status FAILED` |
| `workflow search` | Search executions | None | `--workflow`, `--status`, `--correlation-id`, `--count`, `--start-time-after`, `--start-time-before`, `--json` | `conductor workflow search --workflow my_workflow --status FAILED` |
| `workflow terminate <id>` | Terminate execution | workflow ID | | `conductor workflow terminate abc-123` |
| `workflow pause <id>` | Pause execution | workflow ID | | `conductor workflow pause abc-123` |
| `workflow resume <id>` | Resume paused execution | workflow ID | | `conductor workflow resume abc-123` |
Expand Down Expand Up @@ -713,16 +713,22 @@ conductor workflow search --workflow my_workflow \
--start-time-after "2025-01-01" \
--start-time-before "2025-01-31"

# Find every execution tagged with a correlation ID at start time
conductor workflow start --workflow my_workflow --correlation "order-12345"
conductor workflow search --correlation-id "order-12345"

# Combine filters
conductor workflow search --workflow my_workflow \
--status RUNNING \
--correlation-id "order-12345" \
--start-time-after "2025-01-01 10:00:00" \
--count 100
```

**Search flags:**
- `--workflow <name>` - Filter by workflow name
- `--status <status>` - Filter by status (COMPLETED, FAILED, RUNNING, PAUSED, TERMINATED, TIMED_OUT)
- `--correlation-id <id>` - Filter by correlation ID, the value given to `workflow start --correlation`. Matched exactly; quote ids containing spaces.
- `--count <n>` - Number of results (max 1000, default 10)
- `--start-time-after <time>` - Started after time (formats: YYYY-MM-DD HH:MM:SS, YYYY-MM-DD, or epoch milliseconds)
- `--start-time-before <time>` - Started before time (same formats)
Expand Down
1 change: 1 addition & 0 deletions README.md
Original file line number Diff line number Diff line change
Expand Up @@ -296,6 +296,7 @@ conductor workflow <command> [arguments] [flags]
**Workflow Search Options:**
- `-w, --workflow` - Filter by workflow name
- `-s, --status` - Filter by status: `COMPLETED`, `FAILED`, `PAUSED`, `RUNNING`, `TERMINATED`, `TIMED_OUT`
- `--correlation-id` - Filter by correlation ID (the value passed to `start --correlation`)
- `-c, --count` - Number of results (default: 10, max: 1000)
- `--start-time-after` - Filter by start time (format: `YYYY-MM-DD HH:MM:SS`, `YYYY-MM-DD`, or epoch ms)
- `--start-time-before` - Filter by start time
Expand Down
70 changes: 69 additions & 1 deletion cmd/execution_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -15,6 +15,7 @@
package cmd

import (
"strconv"
"testing"
"time"
)
Expand Down Expand Up @@ -178,4 +179,71 @@ func TestParseTimeToEpochMillis_EdgeCases(t *testing.T) {
t.Errorf("expected 9223372036854775807, got %d", result)
}
})
}
}
func TestBuildWorkflowSearchQuery(t *testing.T) {
startTime := "2023-12-25 15:30:45"
startTimeMs := strconv.FormatInt(time.Date(2023, 12, 25, 15, 30, 45, 0, time.UTC).Unix()*1000, 10)

tests := []struct {
name string
workflowName string
status string
correlationId string
startTimeAfter string
startTimeBefore string
expected string
expectError bool
}{
{
name: "no filters",
expected: "",
},
{
name: "correlation id only",
correlationId: "order-123",
expected: `correlationId = "order-123"`,
},
{
name: "correlation id with spaces stays one term",
correlationId: "order 123",
expected: `correlationId = "order 123"`,
},
{
name: "correlation id combined with other filters",
workflowName: "my_workflow",
status: "COMPLETED",
correlationId: "order-123",
expected: `workflowType IN (my_workflow) AND status IN (COMPLETED) AND correlationId = "order-123"`,
},
{
name: "correlation id combined with time range",
correlationId: "order-123",
startTimeAfter: startTime,
startTimeBefore: startTime,
expected: `correlationId = "order-123" AND startTime>` + startTimeMs + ` AND startTime<` + startTimeMs,
},
{
name: "invalid start-time-after",
startTimeAfter: "not-a-time",
expectError: true,
},
}

for _, tt := range tests {
t.Run(tt.name, func(t *testing.T) {
query, err := buildWorkflowSearchQuery(tt.workflowName, tt.status, tt.correlationId, tt.startTimeAfter, tt.startTimeBefore)
if tt.expectError {
if err == nil {
t.Errorf("expected an error, got query %q", query)
}
return
}
if err != nil {
t.Errorf("unexpected error: %v", err)
}
if query != tt.expected {
t.Errorf("expected %q, got %q", tt.expected, query)
}
})
}
}
65 changes: 41 additions & 24 deletions cmd/workflow.go
Original file line number Diff line number Diff line change
Expand Up @@ -816,59 +816,75 @@ func debugSearchWorkflows(freeText, query string, count int32) error {
return nil
}

func searchWorkflowExecutions(cmd *cobra.Command, args []string) error {
debug, _ := cmd.Flags().GetBool("debug")

freeText := "*"
if len(args) == 1 {
freeText = args[0]
}

count, _ := cmd.Flags().GetInt32("count")
if count > 1000 {
//fmt.Println("count exceeds max allowed 1000. Will only show the first 1000 matching results")
//count = 1000
} else if count == 0 {
count = 10
}

// Build query dynamically with AND conditions
// buildWorkflowSearchQuery combines the search filters into a single query
// expression. Values that can contain spaces are quoted so the server-side
// query parser keeps them as one term.
func buildWorkflowSearchQuery(workflowName, status, correlationId, startTimeAfter, startTimeBefore string) (string, error) {
var queryParts []string

// Workflow name filter
workflowName, _ := cmd.Flags().GetString("workflow")
if workflowName != "" {
queryParts = append(queryParts, "workflowType IN ("+workflowName+")")
}

// Status filter
status, _ := cmd.Flags().GetString("status")
if status != "" {
queryParts = append(queryParts, "status IN ("+status+")")
}

// Correlation ID filter
if correlationId != "" {
queryParts = append(queryParts, "correlationId = \""+correlationId+"\"")
}

// Start time filter (after)
startTimeAfter, _ := cmd.Flags().GetString("start-time-after")
if startTimeAfter != "" {
startTimeAfterMs, err := parseTimeToEpochMillis(startTimeAfter)
if err != nil {
return fmt.Errorf("invalid start-time-after: %v", err)
return "", fmt.Errorf("invalid start-time-after: %v", err)
}
queryParts = append(queryParts, "startTime>"+strconv.FormatInt(startTimeAfterMs, 10))
}

// Start time filter (before)
startTimeBefore, _ := cmd.Flags().GetString("start-time-before")
if startTimeBefore != "" {
startTimeBeforeMs, err := parseTimeToEpochMillis(startTimeBefore)
if err != nil {
return fmt.Errorf("invalid start-time-before: %v", err)
return "", fmt.Errorf("invalid start-time-before: %v", err)
}
queryParts = append(queryParts, "startTime<"+strconv.FormatInt(startTimeBeforeMs, 10))
}

// Combine all query parts with AND
query := strings.Join(queryParts, " AND ")
return strings.Join(queryParts, " AND "), nil
}

func searchWorkflowExecutions(cmd *cobra.Command, args []string) error {
debug, _ := cmd.Flags().GetBool("debug")

freeText := "*"
if len(args) == 1 {
freeText = args[0]
}

count, _ := cmd.Flags().GetInt32("count")
if count > 1000 {
//fmt.Println("count exceeds max allowed 1000. Will only show the first 1000 matching results")
//count = 1000
} else if count == 0 {
count = 10
}

workflowName, _ := cmd.Flags().GetString("workflow")
status, _ := cmd.Flags().GetString("status")
correlationId, _ := cmd.Flags().GetString("correlation-id")
startTimeAfter, _ := cmd.Flags().GetString("start-time-after")
startTimeBefore, _ := cmd.Flags().GetString("start-time-before")

query, err := buildWorkflowSearchQuery(workflowName, status, correlationId, startTimeAfter, startTimeBefore)
if err != nil {
return err
}

// Debug mode: make raw HTTP request to see server response
if debug {
Expand Down Expand Up @@ -1455,6 +1471,7 @@ func init() {
searchExecutionCmd.Flags().Int32P("count", "c", 10, "No of workflow executions to return (max 1000)")
searchExecutionCmd.Flags().StringP("status", "s", "", "Filter by status one of (COMPLETED, FAILED, PAUSED, RUNNING, TERMINATED, TIMED_OUT)")
searchExecutionCmd.Flags().StringP("workflow", "w", "", "Workflow name")
searchExecutionCmd.Flags().String("correlation-id", "", "Filter by correlation ID (set with 'workflow start --correlation')")
searchExecutionCmd.Flags().String("start-time-after", "", "Filter executions started after this time (YYYY-MM-DD HH:MM:SS, YYYY-MM-DD, or epoch ms)")
searchExecutionCmd.Flags().String("start-time-before", "", "Filter executions started before this time (YYYY-MM-DD HH:MM:SS, YYYY-MM-DD, or epoch ms)")
searchExecutionCmd.Flags().Bool("json", false, "Output as JSON")
Expand Down
42 changes: 37 additions & 5 deletions test/e2e/search.bats
Original file line number Diff line number Diff line change
Expand Up @@ -7,6 +7,7 @@

WORKFLOW_NAME="cli_e2e_test_workflow_2"
WORKFLOW_FILE="test/e2e/test-workflow-2.json"
CORRELATION_ID="cli e2e search correlation"

setup() {
if [ ! -f "./conductor" ]; then
Expand All @@ -31,7 +32,7 @@ get_workflow_id() {
}

@test "2. Start workflow execution" {
run bash -c "./conductor workflow start --workflow '$WORKFLOW_NAME' 2>/dev/null"
run bash -c "./conductor workflow start --workflow '$WORKFLOW_NAME' --correlation '$CORRELATION_ID' 2>/dev/null"
echo "Output: $output"
[ "$status" -eq 0 ]

Expand Down Expand Up @@ -113,7 +114,38 @@ get_workflow_id() {
echo "$output" | grep -q "$WORKFLOW_NAME"
}

@test "10. Search for non-existent workflow returns no results" {
@test "10. Search by correlation ID" {
WORKFLOW_ID=$(cat /tmp/search_workflow_id.txt)
[ -n "$WORKFLOW_ID" ]

run bash -c "./conductor workflow search --correlation-id '$CORRELATION_ID' 2>/dev/null"
echo "Output: $output"
[ "$status" -eq 0 ]
echo "$output" | grep -q "$WORKFLOW_ID"
}

@test "11. Search by correlation ID combined with other filters" {
WORKFLOW_ID=$(cat /tmp/search_workflow_id.txt)
[ -n "$WORKFLOW_ID" ]

run bash -c "./conductor workflow search --workflow '$WORKFLOW_NAME' --correlation-id '$CORRELATION_ID' --status COMPLETED 2>/dev/null"
echo "Output: $output"
[ "$status" -eq 0 ]
echo "$output" | grep -q "$WORKFLOW_ID"
}

@test "12. Search by unknown correlation ID returns no results" {
WORKFLOW_ID=$(cat /tmp/search_workflow_id.txt)
[ -n "$WORKFLOW_ID" ]

run bash -c "./conductor workflow search --correlation-id 'nonexistent_correlation_xyz_12345' 2>/dev/null"
echo "Output: $output"
[ "$status" -eq 0 ]
RESULT_COUNT=$(echo "$output" | grep -c "$WORKFLOW_ID" || true)
[ "$RESULT_COUNT" -eq 0 ]
}

@test "13. Search for non-existent workflow returns no results" {
run bash -c "./conductor workflow search --workflow 'nonexistent_workflow_xyz_12345' 2>/dev/null"
echo "Output: $output"
[ "$status" -eq 0 ]
Expand All @@ -123,7 +155,7 @@ get_workflow_id() {
[ "$RESULT_COUNT" -eq 0 ]
}

@test "11. Search by status RUNNING returns no matches for completed workflow" {
@test "14. Search by status RUNNING returns no matches for completed workflow" {
run bash -c "./conductor workflow search --workflow '$WORKFLOW_NAME' --status RUNNING --count 1 2>/dev/null"
echo "Output: $output"
[ "$status" -eq 0 ]
Expand All @@ -135,7 +167,7 @@ get_workflow_id() {
# Cleanup
# ==========================================

@test "12. Cleanup - Delete workflow execution" {
@test "15. Cleanup - Delete workflow execution" {
WORKFLOW_ID=$(cat /tmp/search_workflow_id.txt)
[ -n "$WORKFLOW_ID" ]

Expand All @@ -145,7 +177,7 @@ get_workflow_id() {
[[ "$output" == *"deleted successfully"* ]]
}

@test "13. Cleanup - Delete workflow definition" {
@test "16. Cleanup - Delete workflow definition" {
run bash -c "./conductor workflow delete '$WORKFLOW_NAME' 1 -y 2>/dev/null"
echo "Output: $output"
[ "$status" -eq 0 ]
Expand Down
Loading