Class: Hatchet::AdminClient

Inherits:
Object
  • Object
show all
Defined in:
lib/hatchet-sdk.rb,
sig/hatchet-sdk.rbs

Overview

Admin client for triggering and scheduling workflows.

Delegates to the gRPC admin client for actual RPC calls, while handling context variable propagation for parent-child workflow linking.

Instance Method Summary collapse

Constructor Details

#initialize(client:) ⇒ AdminClient



370
371
372
373
# File 'lib/hatchet-sdk.rb', line 370

def initialize(client:)
  @client = client
  @spawn_indices = ContextVars::SpawnIndexTracker.new
end

Instance Method Details

#schedule_workflow(workflow, time, input: {}, options: nil) ⇒ Object

Schedule a workflow for future execution.



468
469
470
471
472
# File 'lib/hatchet-sdk.rb', line 468

def schedule_workflow(workflow, time, input: {}, options: nil)
  name = workflow.respond_to?(:name) ? workflow.name : workflow.to_s
  opts = build_trigger_options(options)
  @client.admin_grpc.schedule_workflow(name, run_at: time, input: input, options: opts)
end

#trigger_workflow(workflow_or_task, input, options: nil) ⇒ Hash

Trigger a workflow run and wait for result.



381
382
383
384
# File 'lib/hatchet-sdk.rb', line 381

def trigger_workflow(workflow_or_task, input, options: nil)
  ref = trigger_workflow_no_wait(workflow_or_task, input, options: options)
  ref.result
end

#trigger_workflow_many(workflow_or_task, items, return_exceptions: false) ⇒ Array

Trigger many workflow runs and wait for all results.



412
413
414
415
416
417
418
419
420
421
422
423
424
425
426
427
428
429
430
431
432
# File 'lib/hatchet-sdk.rb', line 412

def trigger_workflow_many(workflow_or_task, items, return_exceptions: false)
  refs = trigger_workflow_many_no_wait(workflow_or_task, items)

  # Collect results concurrently using threads so that all subscriptions
  # are sent at once rather than serially waiting for each one.
  threads = refs.map do |ref|
    Thread.new do
      if return_exceptions
        begin
          ref.result
        rescue StandardError => e
          e
        end
      else
        ref.result
      end
    end
  end

  threads.map(&:value)
end

#trigger_workflow_many_no_wait(workflow_or_task, items) ⇒ Array<WorkflowRunRef>

Trigger many workflow runs without waiting.

Uses bulk gRPC triggering for efficiency (batched by 1000).



441
442
443
444
445
446
447
448
449
450
451
452
453
454
455
456
457
458
459
# File 'lib/hatchet-sdk.rb', line 441

def trigger_workflow_many_no_wait(workflow_or_task, items)
  name = workflow_or_task.respond_to?(:name) ? workflow_or_task.name : workflow_or_task.to_s

  # Build trigger items with context vars for parent-child linking
  trigger_items = items.map do |item|
    input = item[:input] || {}
    opts = build_trigger_options(item[:options])
    { input: input, options: opts }
  end

  run_ids = @client.admin_grpc.bulk_trigger_workflow(name, trigger_items)
  run_ids.map do |run_id|
    WorkflowRunRef.new(
      workflow_run_id: run_id,
      client: @client,
      listener: @client.workflow_run_listener,
    )
  end
end

#trigger_workflow_no_wait(workflow_or_task, input, options: nil) ⇒ WorkflowRunRef

Trigger a workflow run without waiting for the result.



392
393
394
395
396
397
398
399
400
401
402
403
404
# File 'lib/hatchet-sdk.rb', line 392

def trigger_workflow_no_wait(workflow_or_task, input, options: nil)
  name = workflow_or_task.respond_to?(:name) ? workflow_or_task.name : workflow_or_task.to_s

  # Merge user options with context vars for parent-child linking
  opts = build_trigger_options(options)

  run_id = @client.admin_grpc.trigger_workflow(name, input: input, options: opts)
  WorkflowRunRef.new(
    workflow_run_id: run_id,
    client: @client,
    listener: @client.workflow_run_listener,
  )
end