Skip to content
Closed
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
36 changes: 36 additions & 0 deletions framework/clclient/client.go
Original file line number Diff line number Diff line change
Expand Up @@ -9,6 +9,7 @@ import (
"net/http"
"os"
"regexp"
"strconv"
"strings"
"sync"
"time"
Expand Down Expand Up @@ -1431,3 +1432,38 @@ func ImportP2PKeys(cl []*ChainlinkClient, keys [][]byte) error {
}
return eg.Wait()
}

func (c *ChainlinkClient) ReadWorkflowEvents(workflowID string, sequence int64, limit int) (*WorkflowDebugEvents, *http.Response, error) {
specObj := &WorkflowDebugEvents{}
framework.L.Info().Str(NodeURL, c.Config.URL).Str("ID", workflowID).Int64("sequence", sequence).Int("limit", limit).Msg("Reading Workflow Events")
resp, err := c.APIClient.R().
SetResult(&specObj).
SetPathParams(map[string]string{
"id": workflowID,
}).
SetQueryParams(map[string]string{
"sequence": strconv.FormatInt(sequence, 10),
"limit": strconv.Itoa(limit),
}).
Get("/v2/debug/workflow/{id}/events?sequence={sequence}&limit={limit}")
if err != nil {
return nil, nil, err
}
return specObj, resp.RawResponse, err
}

func (c *ChainlinkClient) ReadOrphanEvents(sequence int64, limit int) (*WorkflowOrphanEvents, *http.Response, error) {
specObj := &WorkflowOrphanEvents{}
framework.L.Info().Str(NodeURL, c.Config.URL).Int64("sequence", sequence).Int("limit", limit).Msg("Reading Workflow Orphan Events")
resp, err := c.APIClient.R().
SetResult(&specObj).
SetQueryParams(map[string]string{
"sequence": strconv.FormatInt(sequence, 10),
"limit": strconv.Itoa(limit),
}).
Get("/v2/debug/workflow/orphan_events?sequence={sequence}&limit={limit}")
if err != nil {
return nil, nil, err
}
return specObj, resp.RawResponse, err
}
42 changes: 42 additions & 0 deletions framework/clclient/models.go
Original file line number Diff line number Diff line change
Expand Up @@ -1419,3 +1419,45 @@ type ForwarderAttributes struct {
CreatedAt time.Time `json:"createdAt"`
UpdatedAt time.Time `json:"updatedAt"`
}

type WorkflowDebugEvents struct {
Data WorkflowDebugEventsData `json:"data"`
}

type WorkflowDebugEventsData struct {
Type string `json:"type"`
ID string `json:"id"`
Attributes WorkflowDebugEventsAttributes `json:"attributes"`
}

type WorkflowDebugEventsAttributes struct {
Events []WorkflowDebugEvent `json:"events"`
}

type WorkflowDebugEvent struct {
Timestamp time.Time `json:"timestamp"`
Sequence int64 `json:"sequence"`
Message []byte `json:"message"` // protobuf encoded
Type string `json:"type"` // protobuf type name
}

type WorkflowOrphanEvents struct {
Data WorkflowOrphanEventsData `json:"data"`
}

type WorkflowOrphanEventsData struct {
Type string `json:"type"`
ID string `json:"id"`
Attributes WorkflowDebugEventsAttributes `json:"attributes"`
}

type WorkflowOrphanEventsAttributes struct {
Events []WorkflowOrphanEvent `json:"events"`
}

type WorkflowOrphanEvent struct {
Timestamp time.Time `json:"timestamp"`
Sequence int64 `json:"sequence"`
Message []byte `json:"message"` // protobuf encoded
Type string `json:"type"` // protobuf type name
}
Loading