Documentation
¶
Overview ¶
Package spooled provides the official Go SDK for Spooled Cloud.
Package spooled provides the official Go SDK for Spooled Cloud.
Spooled Cloud is a modern, scalable job queue and task scheduler service. This SDK provides a complete interface for interacting with all Spooled API endpoints, including job management, queue operations, worker registration, real-time events, scheduling, and workflows.
Basic Usage ¶
Create a client with your API key:
client, err := spooled.NewClient(
spooled.WithAPIKey("sp_live_your_api_key"),
)
if err != nil {
log.Fatal(err)
}
defer client.Close()
Creating Jobs ¶
ctx := context.Background()
resp, err := client.Jobs().Create(ctx, &types.CreateJobRequest{
QueueName: "emails",
Payload: map[string]any{
"to": "user@example.com",
"subject": "Hello!",
},
})
if err != nil {
log.Fatal(err)
}
fmt.Printf("Created job: %s\n", resp.ID)
Configuration Options ¶
The client supports various configuration options:
client, err := spooled.NewClient(
spooled.WithAPIKey("sp_live_..."),
spooled.WithBaseURL("https://api.spooled.cloud"),
spooled.WithTimeout(30 * time.Second),
spooled.WithRetry(spooled.RetryConfig{
MaxRetries: 3,
BaseDelay: time.Second,
}),
spooled.WithDebug(true),
)
Error Handling ¶
All SDK errors implement the error interface and can be inspected:
job, err := client.Jobs().Get(ctx, "job-id")
if err != nil {
if spooled.IsNotFoundError(err) {
fmt.Println("Job not found")
} else if spooled.IsRateLimitError(err) {
var rateLimitErr *spooled.RateLimitError
if errors.As(err, &rateLimitErr) {
fmt.Printf("Rate limited, retry after %d seconds\n", rateLimitErr.GetRetryAfter())
}
} else {
log.Fatal(err)
}
}
Resources ¶
The client provides access to various resources:
- Jobs: Create, list, cancel, retry jobs and manage the dead letter queue
- Queues: List, configure, pause/resume queues
- Workers: Register, heartbeat, and deregister workers
- Schedules: Create and manage scheduled jobs
- Workflows: Create and manage job workflows with dependencies
- Webhooks: Configure outgoing webhooks for events
- Organizations: Manage organizations and usage
- API Keys: Manage API keys
- Billing: Access billing status and portal
- Auth: Authentication and token management
Package spooled provides top-level convenience functions for common operations.
These functions provide shortcuts for the most common SDK operations.
Index ¶
- Constants
- Variables
- func CancelJob(client *Client, jobID string) error
- func CreateJob(client *Client, queueName string, payload map[string]any) (string, error)
- func CreateJobWithOptions(client *Client, req *resources.CreateJobRequest) (string, error)
- func CreateScheduledJob(client *Client, queueName string, payload map[string]any, ...) (string, error)
- func CreateWebhook(client *Client, req *resources.CreateOutgoingWebhookRequest) (*resources.OutgoingWebhook, error)
- func GetDashboard(client *Client) (*resources.DashboardData, error)
- func GetHealth(client *Client) (*resources.HealthResponse, error)
- func GetJob(client *Client, jobID string) (*resources.Job, error)
- func GetQueueStats(client *Client, queueName string) (*resources.QueueStats, error)
- func IsAuthenticationError(err error) bool
- func IsAuthorizationError(err error) bool
- func IsConflictError(err error) bool
- func IsNotFoundError(err error) bool
- func IsPayloadTooLargeError(err error) bool
- func IsRateLimitError(err error) bool
- func IsRetryable(err error) bool
- func IsServerError(err error) bool
- func IsSpooledError(err error) bool
- func IsValidationError(err error) bool
- func ListJobs(client *Client, params *resources.ListJobsParams) ([]resources.Job, error)
- func ListWorkers(client *Client) ([]resources.Worker, error)
- func RegisterWorker(client *Client, req *resources.RegisterWorkerRequest) (*resources.RegisterWorkerResponse, error)
- func ValidateAPIKey(key string) error
- type APIError
- type AuthenticationError
- type AuthorizationError
- type CircuitBreakerConfig
- type CircuitBreakerOpenError
- type Client
- func (c *Client) APIKeys() *resources.APIKeysResource
- func (c *Client) Admin() *resources.AdminResource
- func (c *Client) Auth() *resources.AuthResource
- func (c *Client) Billing() *resources.BillingResource
- func (c *Client) Close() error
- func (c *Client) Dashboard() *resources.DashboardResource
- func (c *Client) GRPC() (*grpc.Client, error)
- func (c *Client) GetConfig() Config
- func (c *Client) Health() *resources.HealthResource
- func (c *Client) Ingest() *resources.IngestResource
- func (c *Client) Jobs() *resources.JobsResource
- func (c *Client) Metrics() *resources.MetricsResource
- func (c *Client) Organizations() *resources.OrganizationsResource
- func (c *Client) Queues() *resources.QueuesResource
- func (c *Client) Realtime() realtime.RealtimeClient
- func (c *Client) RealtimeSSE() *realtime.SSEClient
- func (c *Client) Schedules() *resources.SchedulesResource
- func (c *Client) Webhooks() *resources.WebhooksResource
- func (c *Client) Workers() *resources.WorkersResource
- func (c *Client) Workflows() *resources.WorkflowsResource
- type Config
- type ConflictError
- type Logger
- type LoggerFunc
- type NetworkError
- type NotFoundError
- type Option
- func WithAPIKey(key string) Option
- func WithAccessToken(token string) Option
- func WithAdminKey(key string) Option
- func WithAutoRefreshToken(enabled bool) Option
- func WithBaseURL(url string) Option
- func WithCircuitBreaker(cfg CircuitBreakerConfig) Option
- func WithDebug(enabled bool) Option
- func WithGRPCAddress(addr string) Option
- func WithHeaders(headers map[string]string) Option
- func WithLogger(l Logger) Option
- func WithRefreshToken(token string) Option
- func WithRetry(cfg RetryConfig) Option
- func WithTimeout(d time.Duration) Option
- func WithUserAgent(ua string) Option
- func WithWSURL(url string) Option
- type PayloadTooLargeError
- type RateLimitError
- type RetryConfig
- type ServerError
- type SpooledWorker
- type SpooledWorkerOptions
- type TimeoutError
- type ValidationError
Constants ¶
const ( DefaultBaseURL = "https://api.spooled.cloud" DefaultWSURL = "wss://api.spooled.cloud" DefaultGRPCAddress = "grpc.spooled.cloud:443" DefaultTimeout = 30 * time.Second DefaultAPIVersion = "v1" DefaultAPIBasePath = "/api/v1" )
Default configuration values
Variables ¶
var ( // ErrNoAuth is returned when no authentication is configured. ErrNoAuth = errors.New("no authentication configured: set APIKey or AccessToken") // ErrInvalidAPIKey is returned when the API key format is invalid. ErrInvalidAPIKey = errors.New("invalid API key format: must start with sk_live_, sk_test_, sp_live_, or sp_test_") // ErrCircuitOpen is returned when the circuit breaker is open. ErrCircuitOpen = errors.New("circuit breaker is open") )
Sentinel errors for common conditions
Functions ¶
func CancelJob ¶
CancelJob cancels a job by ID.
Example:
err := spooled.CancelJob(client, "job_123")
if err != nil {
log.Fatal(err)
}
func CreateJob ¶
CreateJob creates a job and returns the job ID.
This is a convenience function for simple job creation.
Example:
jobID, err := spooled.CreateJob(client, "emails", map[string]any{
"to": "user@example.com",
"subject": "Hello!",
})
if err != nil {
log.Fatal(err)
}
func CreateJobWithOptions ¶
func CreateJobWithOptions(client *Client, req *resources.CreateJobRequest) (string, error)
CreateJobWithOptions creates a job with full options.
Example:
jobID, err := spooled.CreateJobWithOptions(client, &resources.CreateJobRequest{
QueueName: "emails",
Payload: map[string]any{"to": "user@example.com"},
Priority: ptr(10),
ScheduledAt: &futureTime,
})
func CreateScheduledJob ¶
func CreateScheduledJob(client *Client, queueName string, payload map[string]any, scheduledAt time.Time) (string, error)
CreateScheduledJob creates a job scheduled for the future.
Example:
future := time.Now().Add(1 * time.Hour)
jobID, err := spooled.CreateScheduledJob(client, "emails", map[string]any{
"to": "user@example.com",
}, future)
func CreateWebhook ¶
func CreateWebhook(client *Client, req *resources.CreateOutgoingWebhookRequest) (*resources.OutgoingWebhook, error)
CreateWebhook creates an outgoing webhook.
Example:
webhook, err := spooled.CreateWebhook(client, &resources.CreateOutgoingWebhookRequest{
Name: "email-events",
URL: "https://api.example.com/webhooks",
Events: []resources.WebhookEvent{resources.WebhookEventJobCompleted},
})
if err != nil {
log.Fatal(err)
}
func GetDashboard ¶
func GetDashboard(client *Client) (*resources.DashboardData, error)
GetDashboard retrieves dashboard data.
Example:
dashboard, err := spooled.GetDashboard(client)
if err != nil {
log.Fatal(err)
}
fmt.Printf("Total jobs: %d\n", dashboard.Jobs.Total)
func GetHealth ¶
func GetHealth(client *Client) (*resources.HealthResponse, error)
GetHealth checks the health of the Spooled service.
Example:
health, err := spooled.GetHealth(client)
if err != nil {
log.Fatal(err)
}
fmt.Printf("Status: %s\n", health.Status)
func GetJob ¶
GetJob retrieves a job by ID.
Example:
job, err := spooled.GetJob(client, "job_123")
if err != nil {
log.Fatal(err)
}
fmt.Printf("Job status: %s\n", job.Status)
func GetQueueStats ¶
func GetQueueStats(client *Client, queueName string) (*resources.QueueStats, error)
GetQueueStats retrieves statistics for a queue.
Example:
stats, err := spooled.GetQueueStats(client, "emails")
if err != nil {
log.Fatal(err)
}
fmt.Printf("Pending jobs: %d\n", stats.PendingJobs)
func IsAuthenticationError ¶
IsAuthenticationError returns true if the error is a 401 authentication error.
func IsAuthorizationError ¶ added in v1.0.14
IsAuthorizationError returns true if the error is a 403 authorization error.
func IsConflictError ¶ added in v1.0.14
IsConflictError returns true if the error is a 409 conflict error.
func IsNotFoundError ¶
IsNotFoundError returns true if the error is a 404 not-found error.
func IsPayloadTooLargeError ¶ added in v1.0.14
IsPayloadTooLargeError returns true if the error is a 413 payload-too-large error.
func IsRateLimitError ¶
IsRateLimitError returns true if the error is a 429 rate-limit error.
func IsRetryable ¶
IsRetryable returns true if the error is retryable (5xx, 429, network, or timeout errors).
func IsServerError ¶ added in v1.0.14
IsServerError returns true if the error is a 5xx server error.
func IsSpooledError ¶
IsSpooledError returns true if the error is (or wraps) a Spooled SDK API error.
func IsValidationError ¶
IsValidationError returns true if the error is a 400/422 validation error.
func ListJobs ¶
ListJobs lists jobs with optional filters.
Example:
jobs, err := spooled.ListJobs(client, &resources.ListJobsParams{
QueueName: &queueName,
Status: &status,
Limit: ptr(50),
})
func ListWorkers ¶
ListWorkers lists all workers.
Example:
workers, err := spooled.ListWorkers(client)
if err != nil {
log.Fatal(err)
}
for _, w := range workers {
fmt.Printf("Worker: %s (%s)\n", w.ID, w.Status)
}
func RegisterWorker ¶
func RegisterWorker(client *Client, req *resources.RegisterWorkerRequest) (*resources.RegisterWorkerResponse, error)
RegisterWorker registers a new worker.
Example:
resp, err := spooled.RegisterWorker(client, &resources.RegisterWorkerRequest{
QueueName: "emails",
Hostname: "worker-01",
})
if err != nil {
log.Fatal(err)
}
fmt.Printf("Registered worker: %s\n", resp.ID)
func ValidateAPIKey ¶
ValidateAPIKey validates the format of an API key. Accepts both sk_ (production keys) and sp_ (documentation examples) prefixes.
The key is trimmed of surrounding whitespace (including trailing newlines) before validation so callers reading keys from files or environment variables do not need to normalize input themselves. This mirrors the trimming performed by `resolveConfig` so `ValidateAPIKey` and the eventual `Authorization` header agree on what a "valid" key is.
Types ¶
type APIError ¶
APIError is the base error type for all Spooled SDK errors.
Every typed error below embeds *APIError, so errors.As(err, &apiErr) with an *APIError target succeeds for any Spooled API error. Use AsSpooledError for a convenient accessor.
func AsSpooledError ¶
AsSpooledError attempts to extract the underlying Spooled APIError from err. It returns the *APIError and true if err is (or wraps) a Spooled API error.
type AuthenticationError ¶
type AuthenticationError = httpx.AuthenticationError
AuthenticationError represents a 401 error.
type AuthorizationError ¶
type AuthorizationError = httpx.AuthorizationError
AuthorizationError represents a 403 error.
type CircuitBreakerConfig ¶
type CircuitBreakerConfig struct {
// Enabled determines if circuit breaker is active.
Enabled bool
// FailureThreshold is the number of failures before opening the circuit.
FailureThreshold int
// SuccessThreshold is the number of successes needed to close the circuit.
SuccessThreshold int
// Timeout is the duration the circuit stays open before allowing a test request.
Timeout time.Duration
}
CircuitBreakerConfig configures the circuit breaker.
func DefaultCircuitBreakerConfig ¶
func DefaultCircuitBreakerConfig() CircuitBreakerConfig
DefaultCircuitBreakerConfig returns the default circuit breaker configuration.
type CircuitBreakerOpenError ¶
type CircuitBreakerOpenError = httpx.CircuitBreakerOpenError
CircuitBreakerOpenError is returned when the circuit breaker is open.
type Client ¶
type Client struct {
// contains filtered or unexported fields
}
Client is the main Spooled SDK client.
func CreateClient ¶
CreateClient creates a new Spooled client with default configuration.
This is a convenience function equivalent to:
client, err := spooled.NewClient(spooled.WithAPIKey(apiKey))
Example:
client, err := spooled.CreateClient("sp_live_your_api_key")
if err != nil {
log.Fatal(err)
}
defer client.Close()
func (*Client) APIKeys ¶
func (c *Client) APIKeys() *resources.APIKeysResource
APIKeys returns the API Keys resource.
func (*Client) Admin ¶
func (c *Client) Admin() *resources.AdminResource
Admin returns the Admin resource.
func (*Client) Auth ¶
func (c *Client) Auth() *resources.AuthResource
Auth returns the Auth resource.
func (*Client) Billing ¶
func (c *Client) Billing() *resources.BillingResource
Billing returns the Billing resource.
func (*Client) Dashboard ¶
func (c *Client) Dashboard() *resources.DashboardResource
Dashboard returns the Dashboard resource.
func (*Client) GRPC ¶
GRPC returns the gRPC client for high-performance operations.
Note: This method dials the gRPC server the first time it is called. Callers should handle connection errors (e.g., local dev without gRPC running).
func (*Client) Health ¶
func (c *Client) Health() *resources.HealthResource
Health returns the Health resource.
func (*Client) Ingest ¶
func (c *Client) Ingest() *resources.IngestResource
Ingest returns the Ingest resource.
func (*Client) Jobs ¶
func (c *Client) Jobs() *resources.JobsResource
Jobs returns the Jobs resource.
func (*Client) Metrics ¶
func (c *Client) Metrics() *resources.MetricsResource
Metrics returns the Metrics resource.
func (*Client) Organizations ¶
func (c *Client) Organizations() *resources.OrganizationsResource
Organizations returns the Organizations resource.
func (*Client) Queues ¶
func (c *Client) Queues() *resources.QueuesResource
Queues returns the Queues resource.
func (*Client) Realtime ¶ added in v1.0.11
func (c *Client) Realtime() realtime.RealtimeClient
Realtime returns a realtime client for WebSocket/SSE event streaming.
func (*Client) RealtimeSSE ¶ added in v1.0.11
RealtimeSSE returns an SSE realtime client for one-way event streaming.
func (*Client) Schedules ¶
func (c *Client) Schedules() *resources.SchedulesResource
Schedules returns the Schedules resource.
func (*Client) Webhooks ¶
func (c *Client) Webhooks() *resources.WebhooksResource
Webhooks returns the Webhooks resource.
func (*Client) Workers ¶
func (c *Client) Workers() *resources.WorkersResource
Workers returns the Workers resource.
func (*Client) Workflows ¶
func (c *Client) Workflows() *resources.WorkflowsResource
Workflows returns the Workflows resource.
type Config ¶
type Config struct {
// APIKey is the API key for authentication (keys start with sp_live_ or sp_test_).
APIKey string
// AccessToken is a JWT access token (alternative to API key).
AccessToken string
// RefreshToken is a JWT refresh token for automatic token renewal.
RefreshToken string
// AdminKey is the admin API key (for /api/v1/admin/* endpoints; uses X-Admin-Key header).
AdminKey string
// BaseURL is the base URL for the REST API.
BaseURL string
// WSURL is the WebSocket URL for realtime events.
WSURL string
// GRPCAddress is the gRPC server address.
GRPCAddress string
// Timeout is the request timeout.
Timeout time.Duration
// Retry is the retry configuration.
Retry RetryConfig
// CircuitBreaker is the circuit breaker configuration.
CircuitBreaker CircuitBreakerConfig
// Headers are additional headers to include in all requests.
Headers map[string]string
// UserAgent is the custom user agent string.
UserAgent string
// Logger is the debug logger.
Logger Logger
// AutoRefreshToken enables automatic token refresh.
AutoRefreshToken bool
}
Config holds the SDK configuration.
type LoggerFunc ¶
LoggerFunc is a function adapter for Logger.
func (LoggerFunc) Debug ¶
func (f LoggerFunc) Debug(msg string, keysAndValues ...any)
Debug implements Logger.
type NetworkError ¶
type NetworkError = httpx.NetworkError
NetworkError represents a network-level error. It is always retryable.
type Option ¶
type Option func(*Config)
Option is a functional option for configuring the client.
func WithAPIKey ¶
WithAPIKey sets the API key for authentication.
func WithAccessToken ¶
WithAccessToken sets the JWT access token.
func WithAutoRefreshToken ¶
WithAutoRefreshToken enables or disables automatic token refresh.
func WithBaseURL ¶
WithBaseURL sets the base URL for the REST API.
func WithCircuitBreaker ¶
func WithCircuitBreaker(cfg CircuitBreakerConfig) Option
WithCircuitBreaker sets the circuit breaker configuration.
func WithGRPCAddress ¶
WithGRPCAddress sets the gRPC server address.
func WithHeaders ¶
WithHeaders sets additional headers for all requests.
func WithRefreshToken ¶
WithRefreshToken sets the JWT refresh token.
func WithUserAgent ¶
WithUserAgent sets a custom user agent string.
type PayloadTooLargeError ¶
type PayloadTooLargeError = httpx.PayloadTooLargeError
PayloadTooLargeError represents a 413 error.
type RateLimitError ¶
type RateLimitError = httpx.RateLimitError
RateLimitError represents a 429 error. It carries rate-limit metadata parsed from the response headers (RetryAfter, Limit, Remaining, Reset).
type RetryConfig ¶
type RetryConfig struct {
// MaxRetries is the maximum number of retry attempts.
MaxRetries int
// BaseDelay is the initial delay before the first retry.
BaseDelay time.Duration
// MaxDelay is the maximum delay between retries.
MaxDelay time.Duration
// Factor is the exponential backoff multiplier.
Factor float64
// Jitter enables randomized jitter on retry delays.
Jitter bool
}
RetryConfig configures retry behavior for failed requests.
func DefaultRetryConfig ¶
func DefaultRetryConfig() RetryConfig
DefaultRetryConfig returns the default retry configuration.
type ServerError ¶
type ServerError = httpx.ServerError
ServerError represents a 5xx error. It is always retryable.
type SpooledWorker ¶
type SpooledWorker struct {
// contains filtered or unexported fields
}
SpooledWorker is a high-level worker for processing jobs.
func CreateWorker ¶
func CreateWorker(client *Client, queueName string, handler func(context.Context, *resources.Job) (any, error)) (*SpooledWorker, error)
CreateWorker creates and starts a Spooled worker.
This is a convenience function that creates, configures, and starts a worker.
Example:
worker, err := spooled.CreateWorker(client, "emails", func(ctx context.Context, job *resources.Job) (any, error) {
fmt.Printf("Processing email to: %v\n", job.Payload["to"])
return map[string]any{"sent": true}, nil
})
if err != nil {
log.Fatal(err)
}
defer worker.Stop()
// Worker runs until stopped
select {}
func NewSpooledWorker ¶
func NewSpooledWorker(c *Client, opts SpooledWorkerOptions) *SpooledWorker
NewSpooledWorker creates a new Spooled worker for processing jobs.
Example:
w := spooled.NewSpooledWorker(client, spooled.SpooledWorkerOptions{
QueueName: "emails",
Concurrency: 10,
})
defer w.Stop()
w.Process(func(ctx context.Context, job *resources.Job) (any, error) {
// Process job
return map[string]any{"processed": true}, nil
})
if err := w.Start(); err != nil {
log.Fatal(err)
}
type SpooledWorkerOptions ¶
type SpooledWorkerOptions struct {
// QueueName is the name of the queue to process jobs from.
QueueName string
// Concurrency is the maximum number of jobs to process concurrently (default: 5).
Concurrency int
// PollInterval is how often to poll for new jobs (default: 1s).
PollInterval time.Duration
// LeaseDuration is the job lease duration in seconds (default: 30).
LeaseDuration int
// Hostname is the worker hostname (default: auto-detected).
Hostname string
// WorkerType identifies the type of worker (default: "go").
WorkerType string
// Version is the worker version (default: SDK version).
Version string
// Metadata contains additional worker metadata.
Metadata map[string]string
// contains filtered or unexported fields
}
SpooledWorkerOptions configures a Spooled worker.
type TimeoutError ¶
type TimeoutError = httpx.TimeoutError
TimeoutError represents a timeout error. It is always retryable.
type ValidationError ¶
type ValidationError = httpx.ValidationError
ValidationError represents a 400/422 error.
Directories
¶
| Path | Synopsis |
|---|---|
|
Package grpc provides a gRPC client for the Spooled API.
|
Package grpc provides a gRPC client for the Spooled API. |
|
Package realtime provides SSE and WebSocket clients for real-time Spooled events.
|
Package realtime provides SSE and WebSocket clients for real-time Spooled events. |
|
Package resources provides REST resource implementations for the Spooled API.
|
Package resources provides REST resource implementations for the Spooled API. |
|
Package types contains request and response types for the Spooled API.
|
Package types contains request and response types for the Spooled API. |
|
Package worker provides job processing runtime for Spooled queues.
|
Package worker provides job processing runtime for Spooled queues. |