diff --git a/internal/app/domain_types.go b/internal/app/domain_types.go index c63a1f5..5ecfa6f 100644 --- a/internal/app/domain_types.go +++ b/internal/app/domain_types.go @@ -233,7 +233,7 @@ type GitCredentialResult struct { // PipelineLogControl marks the start or end of a workflow step. type PipelineLogControl struct { Kind string `json:"kind"` - Step int `json:"step"` + Step int64 `json:"step"` Time string `json:"time"` Status string `json:"status,omitempty"` Content string `json:"content"` @@ -243,7 +243,7 @@ type PipelineLogControl struct { // PipelineLogData is one line of workflow output. type PipelineLogData struct { - Step int `json:"step"` + Step int64 `json:"step"` Time string `json:"time"` Stream string `json:"stream"` Content string `json:"content"` diff --git a/spindle/cbor.go b/spindle/cbor.go index ad21081..9038a84 100644 --- a/spindle/cbor.go +++ b/spindle/cbor.go @@ -89,12 +89,12 @@ func mapString(m map[string]any, key string) string { return "" } -func mapInt(m map[string]any, key string) int { +func mapInt(m map[string]any, key string) int64 { switch n := m[key].(type) { case int64: - return int(n) - case int: return n + case int: + return int64(n) } return 0 } diff --git a/spindle/ci_cancel_pipeline.go b/spindle/ci_cancel_pipeline.go index f6aa63b..09ee931 100644 --- a/spindle/ci_cancel_pipeline.go +++ b/spindle/ci_cancel_pipeline.go @@ -3,12 +3,10 @@ package spindle import ( "context" "fmt" - - "github.com/bluesky-social/indigo/atproto/syntax" ) func (c *Client) CancelPipeline(ctx context.Context, input CancelPipelineInput) error { - if err := c.Post(ctx, syntax.NSID("sh.tangled.ci.cancelPipeline"), input, nil); err != nil { + if err := c.Post(ctx, nsidCancelPipeline, input, nil); err != nil { return fmt.Errorf("cancel pipeline %q: %w", input.Pipeline, err) } return nil diff --git a/spindle/ci_get_pipeline.go b/spindle/ci_get_pipeline.go index 5cf69be..cd36d21 100644 --- a/spindle/ci_get_pipeline.go +++ b/spindle/ci_get_pipeline.go @@ -3,13 +3,11 @@ package spindle import ( "context" "fmt" - - "github.com/bluesky-social/indigo/atproto/syntax" ) func (c *Client) GetPipeline(ctx context.Context, pipelineID string) (*Pipeline, error) { var pipeline Pipeline - if err := c.Get(ctx, syntax.NSID("sh.tangled.ci.getPipeline"), map[string]any{"pipeline": pipelineID}, &pipeline); err != nil { + if err := c.Get(ctx, nsidGetPipeline, map[string]any{"pipeline": pipelineID}, &pipeline); err != nil { return nil, fmt.Errorf("get pipeline %q: %w", pipelineID, err) } return &pipeline, nil diff --git a/spindle/ci_query_pipelines.go b/spindle/ci_query_pipelines.go index 83a2bdc..35443f5 100644 --- a/spindle/ci_query_pipelines.go +++ b/spindle/ci_query_pipelines.go @@ -3,8 +3,6 @@ package spindle import ( "context" "fmt" - - "github.com/bluesky-social/indigo/atproto/syntax" ) func (c *Client) QueryPipelines(ctx context.Context, repoDID, cursor string) (*QueryPipelinesOutput, error) { @@ -13,7 +11,7 @@ func (c *Client) QueryPipelines(ctx context.Context, repoDID, cursor string) (*Q params["cursor"] = cursor } var output QueryPipelinesOutput - if err := c.Get(ctx, syntax.NSID("sh.tangled.ci.queryPipelines"), params, &output); err != nil { + if err := c.Get(ctx, nsidQueryPipelines, params, &output); err != nil { return nil, fmt.Errorf("query pipelines for %q: %w", repoDID, err) } return &output, nil @@ -21,7 +19,7 @@ func (c *Client) QueryPipelines(ctx context.Context, repoDID, cursor string) (*Q func (c *Client) QueryLatestPipeline(ctx context.Context, repoDID string) (*QueryPipelinesOutput, error) { var output QueryPipelinesOutput - if err := c.Get(ctx, syntax.NSID("sh.tangled.ci.queryPipelines"), map[string]any{"repo": repoDID, "limit": 1}, &output); err != nil { + if err := c.Get(ctx, nsidQueryPipelines, map[string]any{"repo": repoDID, "limit": 1}, &output); err != nil { return nil, fmt.Errorf("query latest pipeline for %q: %w", repoDID, err) } return &output, nil diff --git a/spindle/ci_subscribe_pipeline_logs.go b/spindle/ci_subscribe_pipeline_logs.go index 4a60856..10188b8 100644 --- a/spindle/ci_subscribe_pipeline_logs.go +++ b/spindle/ci_subscribe_pipeline_logs.go @@ -25,7 +25,7 @@ func (c *Client) SubscribePipelineLogs(ctx context.Context, pipelineID string, w default: return fmt.Errorf("spindle host must be http(s), got %q", u.Scheme) } - u.Path = "/xrpc/sh.tangled.ci.subscribePipelineLogs" + u.Path = "/xrpc/" + nsidSubscribeLogs.String() query := url.Values{"pipeline": []string{pipelineID}} if len(workflows) > 0 { query["workflows"] = workflows @@ -74,6 +74,8 @@ func (c *Client) SubscribePipelineLogs(ctx context.Context, pipelineID string, w } } +// decodeLogEvent decodes one WebSocket frame into an event. Frames with fewer +// than two CBOR maps or an unknown header type are skipped and return nil. func decodeLogEvent(data []byte) (*PipelineLogEvent, error) { maps, err := decodeCBORMaps(data) if err != nil { diff --git a/spindle/ci_trigger_pipeline.go b/spindle/ci_trigger_pipeline.go index f820955..23f93dd 100644 --- a/spindle/ci_trigger_pipeline.go +++ b/spindle/ci_trigger_pipeline.go @@ -3,13 +3,11 @@ package spindle import ( "context" "fmt" - - "github.com/bluesky-social/indigo/atproto/syntax" ) func (c *Client) TriggerPipeline(ctx context.Context, input TriggerPipelineInput) (*TriggerPipelineOutput, error) { var output TriggerPipelineOutput - if err := c.Post(ctx, syntax.NSID("sh.tangled.ci.triggerPipeline"), input, &output); err != nil { + if err := c.Post(ctx, nsidTriggerPipeline, input, &output); err != nil { return nil, fmt.Errorf("trigger pipeline: %w", err) } return &output, nil diff --git a/spindle/ci_types.go b/spindle/ci_types.go index 588d9af..03c64e6 100644 --- a/spindle/ci_types.go +++ b/spindle/ci_types.go @@ -1,5 +1,17 @@ package spindle +import ( + "github.com/bluesky-social/indigo/atproto/syntax" +) + +const ( + nsidQueryPipelines syntax.NSID = "sh.tangled.ci.queryPipelines" + nsidGetPipeline syntax.NSID = "sh.tangled.ci.getPipeline" + nsidCancelPipeline syntax.NSID = "sh.tangled.ci.cancelPipeline" + nsidTriggerPipeline syntax.NSID = "sh.tangled.ci.triggerPipeline" + nsidSubscribeLogs syntax.NSID = "sh.tangled.ci.subscribePipelineLogs" +) + // Workflow is one workflow executed by a pipeline. type Workflow struct { ID string `json:"id"` @@ -52,7 +64,7 @@ type TriggerPipelineOutput struct { // PipelineLogControl marks the start or end of a workflow step. type PipelineLogControl struct { Kind string `json:"kind"` - Step int `json:"step"` + Step int64 `json:"step"` Time string `json:"time"` Status string `json:"status,omitempty"` Content string `json:"content"` @@ -62,7 +74,7 @@ type PipelineLogControl struct { // PipelineLogData is one line of workflow output. type PipelineLogData struct { - Step int `json:"step"` + Step int64 `json:"step"` Time string `json:"time"` Stream string `json:"stream"` Content string `json:"content"`