diff --git a/.env.example b/.env.example
index 3887ddd..52485ad 100644
--- a/.env.example
+++ b/.env.example
@@ -1,9 +1,9 @@
-ORIGIN=https://pipes.yourdomain.com
-PORT=3000
-NODE_ENV=production
-DATABASE_URL=data/pipes.db
+# Pipes Secrets
+# All other configuration is in config.yaml
+# Copy this file to .env and fill in the secrets
-# Indiko OAuth Configuration
-INDIKO_CLIENT_ID=ikc_xxxxxxxxxxxxxxxxxxxxx
-INDIKO_CLIENT_SECRET=iks_xxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxx
-INDIKO_ORIGIN=https://indiko.dunkirk.sh
+# OAuth (Indiko)
+INDIKO_CLIENT_SECRET=your_client_secret_here
+
+# Session (generate with: openssl rand -base64 32)
+SESSION_SECRET=your_random_secret_here
diff --git a/.gitignore b/.gitignore
index c4591a0..5e2dda8 100644
--- a/.gitignore
+++ b/.gitignore
@@ -1,12 +1,23 @@
-node_modules/
-.wrangler/
-dist/
-.dev.vars
+# Go
+pipes
+*.exe
+*.dll
+*.so
+*.dylib
+
+# Environment and config
.env
-bun.lockb
-data/
+config.yaml
+
+# Database
*.db
*.db-shm
*.db-wal
+data/
+# OS
.DS_Store
+
+# IDE
+.idea/
+.vscode/
diff --git a/CLAUDE.md b/CLAUDE.md
new file mode 100644
index 0000000..d961e9a
--- /dev/null
+++ b/CLAUDE.md
@@ -0,0 +1,72 @@
+# Pipes - Project Instructions
+
+This is a Go application following Herald's architecture patterns.
+
+## Tech Stack
+
+- **Language**: Go 1.24+
+- **Database**: SQLite with direct SQL (no ORM)
+- **Auth**: Indiko OAuth 2.0 server
+- **Logging**: charmbracelet/log for structured logging
+- **Frontend**: Go html/template + Vanilla JavaScript
+- **Deployment**: Single static binary
+
+## Development Commands
+
+```bash
+# Build the project
+go build -o pipes .
+
+# Run in development
+./pipes serve -c config.yaml
+
+# Run with hot reload (using a tool like air)
+air
+
+# Initialize config files
+./pipes init
+
+# Run tests
+go test ./...
+```
+
+## Configuration
+
+- **YAML Config** (config.yaml): All non-sensitive configuration
+- **Environment** (.env): Secrets only (INDIKO_CLIENT_SECRET, SESSION_SECRET)
+- YAML supports env var expansion: `${VAR}` syntax
+
+See config.yaml.example and .env.example for templates.
+
+## Architecture
+
+Follow Herald's patterns:
+- Clean separation of concerns (config/, store/, auth/, engine/, nodes/, web/)
+- SQLite with WAL mode
+- Structured logging with charm log
+- Graceful shutdown with signal handling
+- Session-based authentication with Indiko OAuth
+
+## Code Style
+
+- Use `gofmt` for formatting
+- Structured logging: `logger.Info("message", "key", value)`
+- Error wrapping: `fmt.Errorf("context: %w", err)`
+- Context propagation for cancellation
+- Foreign key constraints in SQLite
+
+## Design Aesthetic
+
+**Neo-Brutalism** - Bold, geometric design:
+- Space Grotesk font
+- Hard borders (3-4px solid)
+- Hard box shadows (no blur)
+- High contrast colors
+- Sharp, geometric shapes
+- Color palette:
+ - Primary: #2563eb (blue)
+ - Secondary: #ff6b35 (orange)
+ - Auth/Indiko: #AB4967 (pink)
+ - Dark: #26242b (near-black)
+ - Background: #f5f5f0 (warm off-white)
+ - White: #fff
diff --git a/CRUSH.md b/CRUSH.md
new file mode 100644
index 0000000..94d64fa
--- /dev/null
+++ b/CRUSH.md
@@ -0,0 +1,187 @@
+# Crush Memory - Pipes Project
+
+## User Preferences
+
+- Follow Herald's Go architecture patterns
+- Use direct SQL for all database operations (no ORM)
+- Follow neo-brutalist design aesthetic (blue/orange for app, pink only for auth)
+- Use Indiko for authentication (OAuth 2.0 with PKCE)
+- Run on port 3001 (Indiko runs on 3000)
+- Use charmbracelet/log for structured logging
+
+## Architecture Patterns
+
+### Authentication Flow
+- OAuth 2.0 client that uses Indiko for authentication
+- PKCE flow with code verifier/challenge (required by IndieAuth spec)
+- Auto-registration with Indiko (client ID = app URL)
+- Session-based auth with 30-day cookies
+- Automatic token refresh using refresh tokens
+
+### Database (SQLite)
+- Direct SQL queries (no ORM like Kysely or GORM)
+- SQLite with WAL mode for concurrency
+- Foreign key constraints enabled
+- Schema created automatically on startup
+
+### Project Structure
+```
+pipes/
+├── main.go # Entry point, CLI setup
+├── go.mod # Dependencies
+├── config/
+│ ├── app.go # Config struct & loading (like Herald)
+│ └── validate.go # Config validation
+├── store/
+│ ├── db.go # SQLite setup & schema
+│ ├── users.go # User operations
+│ ├── pipes.go # Pipe CRUD
+│ ├── executions.go # Execution history
+│ └── cache.go # Source cache operations
+├── auth/
+│ ├── oauth.go # OAuth 2.0 client (Indiko)
+│ ├── session.go # Session management
+│ └── middleware.go # Auth middleware
+├── engine/
+│ ├── executor.go # Pipeline execution engine
+│ ├── scheduler.go # Cron-based scheduler (Herald pattern)
+│ └── registry.go # Node type registry
+├── nodes/
+│ ├── node.go # Node interface
+│ ├── sources/
+│ │ ├── rss.go # RSS/Atom source
+│ │ └── http.go # HTTP API source
+│ └── transforms/
+│ ├── filter.go # Filter transform
+│ ├── sort.go # Sort transform
+│ ├── limit.go # Limit transform
+│ ├── merge.go # Merge sources
+│ ├── dedupe.go # Deduplicate items
+│ └── extract.go # Extract/transform fields
+└── web/
+ ├── server.go # HTTP server setup
+ ├── handlers.go # Route handlers
+ ├── api.go # JSON API endpoints
+ └── templates/
+ ├── layout.html # Base layout
+ ├── index.html # Landing page
+ ├── dashboard.html # User dashboard
+ ├── editor.html # Visual editor
+ └── style.css # Frutiger Aero styles
+```
+
+### Database Schema
+- **users**: id, indiko_sub, username, name, email, photo, url, role, created_at, updated_at
+ - `indiko_sub` is the "sub" field from Indiko's userinfo endpoint
+ - `role` is either "user" or "admin" (synced from Indiko)
+- **sessions**: id, user_id, access_token, refresh_token, expires_at, created_at
+ - 30-day sessions with automatic token refresh
+- **pipes**: id, user_id, name, description, config (JSON), is_public, created_at, updated_at
+ - `config` is JSON: {version, nodes[], connections[], settings}
+- **scheduled_jobs**: id, pipe_id, cron_expression, next_run_at, last_run_at, enabled, created_at, updated_at
+- **pipe_executions**: id, pipe_id, status, trigger_type, started_at, completed_at, duration_ms, items_processed, error_message, metadata
+- **execution_logs**: id, execution_id, node_id, level, message, timestamp, metadata
+- **source_cache**: id, pipe_id, node_id, cache_key, data, etag, last_modified, expires_at, created_at
+
+### OAuth Configuration
+- **Client ID**: App's URL (e.g., `http://localhost:3001`)
+- **Client Secret**: Optional, for pre-registered clients (stored in .env)
+- **Scopes**: `profile email`
+- **PKCE**: Required (S256 code challenge method)
+- **Callback URL**: `{ORIGIN}/auth/callback`
+
+## Configuration
+
+### YAML Config (config.yaml)
+Contains all non-sensitive configuration:
+```yaml
+# Server settings
+host: localhost
+port: 3001
+origin: http://localhost:3001
+env: development
+log_level: info # debug, info, warn, error, fatal
+
+# Database
+db_path: pipes.db
+
+# OAuth (Indiko)
+indiko_url: http://localhost:3000
+indiko_client_id: http://localhost:3001
+indiko_client_secret: ${INDIKO_CLIENT_SECRET} # Loaded from .env
+oauth_callback_url: http://localhost:3001/auth/callback
+
+# Session
+session_secret: ${SESSION_SECRET} # Loaded from .env
+session_cookie_name: pipes_session
+```
+
+### Environment Variables (.env)
+Contains **only secrets**:
+```env
+# OAuth (Indiko)
+INDIKO_CLIENT_SECRET=your_client_secret_here
+
+# Session (generate with: openssl rand -base64 32)
+SESSION_SECRET=your_random_secret_here
+```
+
+## Routes
+
+### Public
+- `GET /` - Landing page (redirects to /dashboard if authenticated)
+- `GET /auth/login` - Start OAuth flow
+- `GET /auth/callback` - OAuth callback handler
+
+### Authenticated
+- `GET /dashboard` - User dashboard (requires auth)
+- `GET /pipes/:id/edit` - Visual editor (requires auth)
+- `POST /auth/logout` - End session
+- `GET /api/me` - Get current user info
+- `GET /api/pipes` - List user's pipes
+- `GET /api/pipes/:id` - Get pipe config
+- `POST /api/pipes` - Create pipe
+- `PUT /api/pipes/:id` - Update pipe
+- `DELETE /api/pipes/:id` - Delete pipe
+- `POST /api/pipes/:id/execute` - Execute pipe manually
+- `GET /api/pipes/:id/executions` - Execution history
+- `GET /api/executions/:id/logs` - Execution logs
+- `GET /api/node-types` - Available node types
+
+## Code Style
+
+- Use `gofmt` for formatting
+- Structured logging with charmbracelet/log: `logger.Info("message", "key", value)`
+- Error wrapping: `fmt.Errorf("context: %w", err)`
+- Context propagation for cancellation
+- Session cookies named `pipes_session`
+- Authorization header: `Bearer {token}`
+
+## Design Aesthetic
+
+**Neo-Brutalism** - Bold, geometric design:
+- Space Grotesk font family
+- Hard borders (3-4px solid #26242b)
+- Hard box shadows (6-12px offset, no blur)
+- High contrast, sharp edges
+- Uppercase text, tight letter-spacing
+- Color palette:
+ - **App colors**:
+ - Primary: #2563eb (blue) - used for headings, secondary buttons
+ - Secondary: #ff6b35 (orange) - used for primary buttons, accents
+ - Background: #f5f5f0 (warm off-white)
+ - Dark: #26242b (near-black)
+ - White: #fff
+ - **Auth/Indiko colors** (login, error pages):
+ - Auth primary: #AB4967 (muted pink)
+ - Auth text: #fff (white - for text on pink buttons)
+
+## Commands
+
+```bash
+go build -o pipes . # Build
+./pipes serve # Run server
+./pipes init # Initialize config files
+./pipes help # Show help
+./pipes version # Show version
+```
diff --git a/README.md b/README.md
index 25e08ba..be7c7f5 100644
--- a/README.md
+++ b/README.md
@@ -1,9 +1,27 @@
# Pipes
-This is my interperitation of yahoo pipes from back in the day! It is designed to allow you to string together pipelines of data and do cool stuff!
+This is my interpretation of Yahoo Pipes from back in the day! It is designed to allow you to string together pipelines of data and do cool stuff with a modern Frutiger Aero aesthetic!
The canonical repo for this is hosted on tangled over at [`dunkirk.sh/pipes`](https://tangled.org/@dunkirk.sh/pipes)
+## Features
+
+- 🔐 **Passwordless Authentication** - Uses Indiko for OAuth 2.0 authentication with passkeys
+- 🌊 **Visual Pipeline Builder** - Create data flows with an intuitive drag-and-drop interface
+- ⚡ **Scheduled Execution** - Pipes run automatically on cron schedules
+- 📊 **Data Sources** - RSS/Atom feeds and HTTP/REST APIs
+- 🔄 **Transform Operations** - Filter, sort, limit, merge, dedupe, and extract data
+- 🎨 **Neo-Brutalist Design** - Bold, geometric UI matching Indiko's aesthetic
+- 👥 **Role-based Access** - User and admin roles powered by Indiko
+
+## Tech Stack
+
+- **Language**: Go 1.24+
+- **Database**: SQLite with direct SQL
+- **Auth**: [Indiko](https://github.com/taciturnaxolotl/indiko) OAuth 2.0 server
+- **Frontend**: Go html/template + Vanilla JavaScript
+- **Deployment**: Single static binary
+
## Installation
1. Clone the repository:
@@ -13,51 +31,196 @@ git clone https://github.com/taciturnaxolotl/pipes.git
cd pipes
```
-2. Install dependencies:
+2. Build the binary:
```bash
-bun install
+go build -o pipes .
```
-3. Create a `.env` file:
+3. Initialize configuration:
+
+```bash
+./pipes init
+```
+
+This creates a `config.yaml` file with sample configuration and a `.env.example` file for secrets.
+
+Copy the example and add your secrets:
```bash
cp .env.example .env
+# Edit .env with your actual secrets
```
-Configure the following environment variables:
+Example `.env` file:
```env
-ORIGIN=https://pipes.yourdomain.com
-PORT=3000
-NODE_ENV=production
-DATABASE_URL=data/pipes.db
-
-# Indiko OAuth Configuration
-INDIKO_CLIENT_ID=ikc_xxxxxxxxxxxxxxxxxxxxx
-INDIKO_CLIENT_SECRET=iks_xxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxx
-INDIKO_ORIGIN=https://indiko.dunkirk.sh
-INDIKO_REDIRECT_URI=https://pipes.yourdomain.com/auth/callback
+# Pipes Secrets
+# All other configuration is in config.yaml
+# Copy this file to .env and fill in the secrets
+
+# OAuth (Indiko)
+INDIKO_CLIENT_SECRET=your_client_secret_here
+
+# Session (generate with: openssl rand -base64 32)
+SESSION_SECRET=your_random_secret_here
```
-The database will be automatically created at `./data/pipes.db` on first run.
+The database will be automatically created at `./pipes.db` on first run.
4. Set up Indiko OAuth:
- - Go to your Indiko instance
- - Navigate to Admin → OAuth Clients
- - Create a new client with the redirect URI matching your `INDIKO_REDIRECT_URI`
- - Copy the Client ID and Secret to your `.env` file
+
+Pipes uses auto-registration with Indiko, so you can start using it immediately! The client ID is just your app's URL (`http://localhost:3001`).
+
+For production or to use role-based access control, ask your Indiko admin to pre-register your client with a client secret.
5. Start the server:
```bash
-# Development (with hot reload)
-bun run dev
+./pipes serve -c config.yaml
+```
+
+Or run without specifying a config file (uses environment variables from `.env`):
+
+```bash
+./pipes serve
+```
+
+Visit `http://localhost:3001` and sign in with your Indiko account!
+
+## Configuration
+
+Pipes uses a two-file configuration approach (just like Herald):
+
+### YAML Config File (config.yaml)
+
+Contains all non-sensitive configuration:
+
+```bash
+./pipes init # Creates config.yaml and .env.example
+./pipes serve -c config.yaml
+```
+
+Example `config.yaml`:
+```yaml
+# Server settings
+host: localhost
+port: 3001
+origin: http://localhost:3001
+env: development
+log_level: info # debug, info, warn, error, fatal
+
+# Database
+db_path: pipes.db
+
+# OAuth (Indiko)
+indiko_url: http://localhost:3000
+indiko_client_id: http://localhost:3001
+indiko_client_secret: ${INDIKO_CLIENT_SECRET} # Loaded from .env
+oauth_callback_url: http://localhost:3001/auth/callback
+
+# Session
+session_secret: ${SESSION_SECRET} # Loaded from .env
+session_cookie_name: pipes_session
+```
+
+### Environment Variables (.env)
+
+Contains **only secrets** (never commit this file):
-# Production
-bun run start
+```env
+# OAuth (Indiko)
+INDIKO_CLIENT_SECRET=your_client_secret_here
+
+# Session (generate with: openssl rand -base64 32)
+SESSION_SECRET=your_random_secret_here
+```
+
+YAML config supports environment variable expansion using `${VAR}` syntax. Variables are loaded from `.env` file and can be overridden by system environment variables.
+
+**Configuration precedence:** Environment variables > YAML config > defaults
+
+### Log Levels
+
+Set `LOG_LEVEL` (or `log_level` in YAML) to:
+- `debug` - Verbose output for troubleshooting
+- `info` - Standard operational messages (default)
+- `warn` - Warning messages
+- `error` - Error messages only
+- `fatal` - Fatal errors (exits immediately)
+
+Example structured logging output:
+```
+2026/01/10 10:24:05 INFO starting pipes host=localhost port=3001 db_path=pipes.db
+2026/01/10 10:24:05 INFO user authenticated name="John Doe" email="john@example.com"
+```
+
+## Architecture
+
+Pipes follows Herald's clean architecture patterns:
+
+```
+pipes/
+├── main.go # CLI entry point
+├── config/ # Configuration management
+├── store/ # Database operations
+├── auth/ # OAuth 2.0 client & session management
+├── engine/ # Pipeline executor & scheduler
+├── nodes/ # Node type definitions
+│ ├── sources/ # RSS, HTTP API sources
+│ └── transforms/ # Filter, sort, limit operations
+└── web/ # HTTP server & handlers
+ └── templates/ # HTML templates
+```
+
+## OAuth Flow
+
+1. User clicks "Sign in with Indiko"
+2. Redirect to Indiko authorization endpoint with PKCE
+3. User authenticates with passkey on Indiko
+4. User approves scopes (profile, email)
+5. Indiko redirects back with authorization code
+6. Exchange code for access + refresh tokens
+7. Create/update user in local database
+8. Create session with 30-day cookie
+
+## Pipeline Execution
+
+Pipelines are executed using topological sort (Kahn's algorithm):
+
+1. Parse pipe configuration (nodes + connections)
+2. Build dependency graph
+3. Execute nodes in order, passing data between them
+4. Log execution progress
+5. Store results in database
+
+The scheduler runs every minute, checking for pipes that need to execute based on their cron schedules.
+
+## Available Node Types
+
+**Sources:**
+- RSS Feed - Fetch items from RSS/Atom feeds
+- HTTP API - Fetch JSON data from REST APIs (coming soon)
+
+**Transforms:**
+- Filter - Filter items based on field conditions
+- Sort - Sort items by field values
+- Limit - Limit the number of output items
+- Merge - Combine multiple data sources (coming soon)
+- Dedupe - Remove duplicate items (coming soon)
+- Extract - Transform/extract fields (coming soon)
+
+## Development
+
+Build and run:
+
+```bash
+go build -o pipes .
+./pipes serve
```
+The database schema is automatically created on first run.
+
@@ -69,3 +232,4 @@ bun run start
+
diff --git a/auth/middleware.go b/auth/middleware.go
new file mode 100644
index 0000000..880c73b
--- /dev/null
+++ b/auth/middleware.go
@@ -0,0 +1,31 @@
+package auth
+
+import (
+ "context"
+ "net/http"
+
+ "github.com/kierank/pipes/store"
+)
+
+type contextKey string
+
+const userContextKey contextKey = "user"
+
+func (sm *SessionManager) RequireAuth(next http.HandlerFunc) http.HandlerFunc {
+ return func(w http.ResponseWriter, r *http.Request) {
+ user, err := sm.GetCurrentUser(r)
+ if err != nil || user == nil {
+ http.Redirect(w, r, "/auth/login", http.StatusSeeOther)
+ return
+ }
+
+ // Add user to context
+ ctx := context.WithValue(r.Context(), userContextKey, user)
+ next(w, r.WithContext(ctx))
+ }
+}
+
+func GetUserFromContext(ctx context.Context) *store.User {
+ user, _ := ctx.Value(userContextKey).(*store.User)
+ return user
+}
diff --git a/auth/oauth.go b/auth/oauth.go
new file mode 100644
index 0000000..bc7ddfc
--- /dev/null
+++ b/auth/oauth.go
@@ -0,0 +1,242 @@
+package auth
+
+import (
+ "crypto/rand"
+ "crypto/sha256"
+ "encoding/base64"
+ "encoding/json"
+ "fmt"
+ "io"
+ "net/http"
+ "net/url"
+ "strings"
+ "time"
+
+ "github.com/kierank/pipes/config"
+ "github.com/kierank/pipes/store"
+)
+
+type OAuthClient struct {
+ cfg *config.Config
+ db *store.DB
+ states map[string]*PKCEState // In-memory for MVP; use Redis in production
+}
+
+type PKCEState struct {
+ CodeVerifier string
+ RedirectURI string
+ CreatedAt time.Time
+}
+
+type TokenResponse struct {
+ AccessToken string `json:"access_token"`
+ TokenType string `json:"token_type"`
+ ExpiresIn int `json:"expires_in"`
+ RefreshToken string `json:"refresh_token,omitempty"`
+ Scope string `json:"scope,omitempty"`
+}
+
+type UserInfo struct {
+ Sub string `json:"sub"`
+ Username string `json:"username,omitempty"`
+ Name string `json:"name,omitempty"`
+ Email string `json:"email,omitempty"`
+ Photo string `json:"picture,omitempty"`
+ URL string `json:"profile,omitempty"`
+}
+
+func NewOAuthClient(cfg *config.Config, db *store.DB) *OAuthClient {
+ return &OAuthClient{
+ cfg: cfg,
+ db: db,
+ states: make(map[string]*PKCEState),
+ }
+}
+
+func (c *OAuthClient) GetAuthorizationURL() (string, error) {
+ state, err := generateRandomString(32)
+ if err != nil {
+ return "", fmt.Errorf("generate state: %w", err)
+ }
+
+ codeVerifier, err := generateRandomString(64)
+ if err != nil {
+ return "", fmt.Errorf("generate code verifier: %w", err)
+ }
+
+ codeChallenge := generateCodeChallenge(codeVerifier)
+
+ // Store PKCE state (in-memory for now)
+ c.states[state] = &PKCEState{
+ CodeVerifier: codeVerifier,
+ RedirectURI: c.cfg.OAuthCallbackURL,
+ CreatedAt: time.Now(),
+ }
+
+ // Clean up old states (older than 10 minutes)
+ go c.cleanupStates()
+
+ authURL := fmt.Sprintf("%s/auth/authorize?"+
+ "response_type=code&"+
+ "client_id=%s&"+
+ "redirect_uri=%s&"+
+ "state=%s&"+
+ "code_challenge=%s&"+
+ "code_challenge_method=S256&"+
+ "scope=profile%%20email",
+ c.cfg.IndikoURL,
+ url.QueryEscape(c.cfg.IndikoClientID),
+ url.QueryEscape(c.cfg.OAuthCallbackURL),
+ state,
+ codeChallenge,
+ )
+
+ return authURL, nil
+}
+
+func (c *OAuthClient) HandleCallback(state, code string) (*store.User, *store.Session, error) {
+ // Verify state
+ pkceState, ok := c.states[state]
+ if !ok {
+ return nil, nil, fmt.Errorf("invalid state")
+ }
+
+ delete(c.states, state)
+
+ // Exchange code for token
+ tokenResp, err := c.exchangeCode(code, pkceState.CodeVerifier, pkceState.RedirectURI)
+ if err != nil {
+ return nil, nil, fmt.Errorf("exchange code: %w", err)
+ }
+
+ // Fetch user info
+ userInfo, err := c.fetchUserInfo(tokenResp.AccessToken)
+ if err != nil {
+ return nil, nil, fmt.Errorf("fetch user info: %w", err)
+ }
+
+ // Create or update user
+ user, err := c.db.GetUserByIndikoSub(userInfo.Sub)
+ if err != nil {
+ return nil, nil, fmt.Errorf("get user: %w", err)
+ }
+
+ if user == nil {
+ user, err = c.db.CreateUser(userInfo.Sub, userInfo.Username, userInfo.Name, userInfo.Email, userInfo.Photo, userInfo.URL)
+ if err != nil {
+ return nil, nil, fmt.Errorf("create user: %w", err)
+ }
+ } else {
+ // Update user info
+ user.Username = userInfo.Username
+ user.Name = userInfo.Name
+ user.Email = userInfo.Email
+ user.Photo = userInfo.Photo
+ user.URL = userInfo.URL
+ if err := c.db.UpdateUser(user); err != nil {
+ return nil, nil, fmt.Errorf("update user: %w", err)
+ }
+ }
+
+ // Create session
+ expiresAt := time.Now().Add(30 * 24 * time.Hour).Unix() // 30 days
+ session, err := c.db.CreateSession(user.ID, tokenResp.AccessToken, tokenResp.RefreshToken, expiresAt)
+ if err != nil {
+ return nil, nil, fmt.Errorf("create session: %w", err)
+ }
+
+ return user, session, nil
+}
+
+func (c *OAuthClient) exchangeCode(code, codeVerifier, redirectURI string) (*TokenResponse, error) {
+ data := url.Values{}
+ data.Set("grant_type", "authorization_code")
+ data.Set("code", code)
+ data.Set("redirect_uri", redirectURI)
+ data.Set("client_id", c.cfg.IndikoClientID)
+ data.Set("code_verifier", codeVerifier)
+
+ if c.cfg.IndikoClientSecret != "" {
+ data.Set("client_secret", c.cfg.IndikoClientSecret)
+ }
+
+ tokenURL := fmt.Sprintf("%s/auth/token", c.cfg.IndikoURL)
+
+ // Create request with explicit headers
+ req, err := http.NewRequest("POST", tokenURL, strings.NewReader(data.Encode()))
+ if err != nil {
+ return nil, fmt.Errorf("create request: %w", err)
+ }
+ req.Header.Set("Content-Type", "application/x-www-form-urlencoded")
+
+ client := &http.Client{Timeout: 10 * time.Second}
+ resp, err := client.Do(req)
+ if err != nil {
+ return nil, fmt.Errorf("do request: %w", err)
+ }
+ defer resp.Body.Close()
+
+ if resp.StatusCode != http.StatusOK {
+ body, _ := io.ReadAll(resp.Body)
+ return nil, fmt.Errorf("token request failed (URL: %s): %s - %s", tokenURL, resp.Status, string(body))
+ }
+
+ var tokenResp TokenResponse
+ if err := json.NewDecoder(resp.Body).Decode(&tokenResp); err != nil {
+ return nil, fmt.Errorf("decode response: %w", err)
+ }
+
+ return &tokenResp, nil
+}
+
+func (c *OAuthClient) fetchUserInfo(accessToken string) (*UserInfo, error) {
+ userInfoURL := fmt.Sprintf("%s/userinfo", c.cfg.IndikoURL)
+
+ req, err := http.NewRequest("GET", userInfoURL, nil)
+ if err != nil {
+ return nil, fmt.Errorf("create request: %w", err)
+ }
+
+ req.Header.Set("Authorization", "Bearer "+accessToken)
+
+ client := &http.Client{Timeout: 10 * time.Second}
+ resp, err := client.Do(req)
+ if err != nil {
+ return nil, fmt.Errorf("do request: %w", err)
+ }
+ defer resp.Body.Close()
+
+ if resp.StatusCode != http.StatusOK {
+ body, _ := io.ReadAll(resp.Body)
+ return nil, fmt.Errorf("userinfo request failed: %s - %s", resp.Status, string(body))
+ }
+
+ var userInfo UserInfo
+ if err := json.NewDecoder(resp.Body).Decode(&userInfo); err != nil {
+ return nil, fmt.Errorf("decode response: %w", err)
+ }
+
+ return &userInfo, nil
+}
+
+func (c *OAuthClient) cleanupStates() {
+ cutoff := time.Now().Add(-10 * time.Minute)
+ for state, pkceState := range c.states {
+ if pkceState.CreatedAt.Before(cutoff) {
+ delete(c.states, state)
+ }
+ }
+}
+
+func generateRandomString(length int) (string, error) {
+ bytes := make([]byte, length)
+ if _, err := rand.Read(bytes); err != nil {
+ return "", err
+ }
+ return base64.RawURLEncoding.EncodeToString(bytes), nil
+}
+
+func generateCodeChallenge(verifier string) string {
+ hash := sha256.Sum256([]byte(verifier))
+ return base64.RawURLEncoding.EncodeToString(hash[:])
+}
diff --git a/auth/session.go b/auth/session.go
new file mode 100644
index 0000000..650afb0
--- /dev/null
+++ b/auth/session.go
@@ -0,0 +1,77 @@
+package auth
+
+import (
+ "net/http"
+
+ "github.com/gorilla/sessions"
+ "github.com/kierank/pipes/config"
+ "github.com/kierank/pipes/store"
+)
+
+type SessionManager struct {
+ store *sessions.CookieStore
+ db *store.DB
+ cfg *config.Config
+}
+
+func NewSessionManager(cfg *config.Config, db *store.DB) *SessionManager {
+ store := sessions.NewCookieStore([]byte(cfg.SessionSecret))
+ store.Options = &sessions.Options{
+ Path: "/",
+ MaxAge: 30 * 24 * 60 * 60, // 30 days
+ HttpOnly: true,
+ SameSite: http.SameSiteLaxMode,
+ Secure: cfg.Env == "production",
+ }
+
+ return &SessionManager{
+ store: store,
+ db: db,
+ cfg: cfg,
+ }
+}
+
+func (sm *SessionManager) SetSession(w http.ResponseWriter, r *http.Request, sessionID string) error {
+ session, _ := sm.store.Get(r, sm.cfg.SessionCookieName)
+ session.Values["session_id"] = sessionID
+ return session.Save(r, w)
+}
+
+func (sm *SessionManager) GetSessionID(r *http.Request) (string, error) {
+ session, err := sm.store.Get(r, sm.cfg.SessionCookieName)
+ if err != nil {
+ return "", err
+ }
+
+ sessionID, ok := session.Values["session_id"].(string)
+ if !ok {
+ return "", nil
+ }
+
+ return sessionID, nil
+}
+
+func (sm *SessionManager) ClearSession(w http.ResponseWriter, r *http.Request) error {
+ session, _ := sm.store.Get(r, sm.cfg.SessionCookieName)
+ session.Options.MaxAge = -1
+ return session.Save(r, w)
+}
+
+func (sm *SessionManager) GetCurrentUser(r *http.Request) (*store.User, error) {
+ sessionID, err := sm.GetSessionID(r)
+ if err != nil || sessionID == "" {
+ return nil, nil
+ }
+
+ session, err := sm.db.GetSessionByID(sessionID)
+ if err != nil || session == nil {
+ return nil, err
+ }
+
+ user, err := sm.db.GetUserByID(session.UserID)
+ if err != nil {
+ return nil, err
+ }
+
+ return user, nil
+}
diff --git a/config.yaml.example b/config.yaml.example
new file mode 100644
index 0000000..fd14686
--- /dev/null
+++ b/config.yaml.example
@@ -0,0 +1,23 @@
+# Pipes Configuration
+# Secrets are loaded from .env file (see .env.example)
+# Environment variables override values in this file
+
+# Server settings
+host: localhost
+port: 3001
+origin: http://localhost:3001
+env: development
+log_level: info # debug, info, warn, error, fatal
+
+# Database
+db_path: pipes.db
+
+# OAuth (Indiko)
+indiko_url: https://indiko.example.com # Use HTTPS (HTTP redirects cause issues)
+indiko_client_id: http://localhost:3001
+indiko_client_secret: ${INDIKO_CLIENT_SECRET} # Loaded from .env
+oauth_callback_url: http://localhost:3001/auth/callback
+
+# Session
+session_secret: ${SESSION_SECRET} # Loaded from .env
+session_cookie_name: pipes_session
diff --git a/config/app.go b/config/app.go
new file mode 100644
index 0000000..0e2dbea
--- /dev/null
+++ b/config/app.go
@@ -0,0 +1,161 @@
+package config
+
+import (
+ "fmt"
+ "os"
+ "path/filepath"
+ "strconv"
+
+ "github.com/joho/godotenv"
+ "gopkg.in/yaml.v3"
+)
+
+type Config struct {
+ // Server
+ Origin string `yaml:"origin"`
+ Host string `yaml:"host"`
+ Port int `yaml:"port"`
+ Env string `yaml:"env"`
+ LogLevel string `yaml:"log_level"`
+
+ // Database
+ DatabasePath string `yaml:"db_path"`
+
+ // OAuth (Indiko)
+ IndikoURL string `yaml:"indiko_url"`
+ IndikoClientID string `yaml:"indiko_client_id"`
+ IndikoClientSecret string `yaml:"indiko_client_secret"`
+ OAuthCallbackURL string `yaml:"oauth_callback_url"`
+
+ // Session
+ SessionSecret string `yaml:"session_secret"`
+ SessionCookieName string `yaml:"session_cookie_name"`
+}
+
+// Default returns a Config with sensible defaults
+func Default() *Config {
+ return &Config{
+ Origin: "http://localhost:3001",
+ Host: "localhost",
+ Port: 3001,
+ Env: "development",
+ LogLevel: "info",
+ DatabasePath: "pipes.db",
+ IndikoURL: "http://localhost:3000",
+ OAuthCallbackURL: "http://localhost:3001/auth/callback",
+ SessionCookieName: "pipes_session",
+ }
+}
+
+// Load loads configuration from YAML file (if provided) and environment variables
+func Load(path string) (*Config, error) {
+ cfg := Default()
+
+ // Load .env file if it exists (silently ignore if not found)
+ if envPath := findEnvFile(path); envPath != "" {
+ _ = godotenv.Load(envPath)
+ }
+
+ // Load from YAML config file if provided
+ if path != "" {
+ data, err := os.ReadFile(path)
+ if err != nil {
+ return nil, fmt.Errorf("failed to read config file: %w", err)
+ }
+
+ // Expand environment variables in YAML (e.g., ${DATABASE_PATH})
+ expanded := os.Expand(string(data), func(key string) string {
+ return os.Getenv(key)
+ })
+
+ if err := yaml.Unmarshal([]byte(expanded), cfg); err != nil {
+ return nil, fmt.Errorf("failed to parse config file: %w", err)
+ }
+ }
+
+ // Apply environment variable overrides
+ applyEnvOverrides(cfg)
+
+ if err := cfg.Validate(); err != nil {
+ return nil, err
+ }
+
+ return cfg, nil
+}
+
+func (c *Config) Validate() error {
+ if c.SessionSecret == "" {
+ return fmt.Errorf("session_secret is required (set SESSION_SECRET env var)")
+ }
+
+ if c.IndikoClientID == "" {
+ return fmt.Errorf("indiko_client_id is required (set INDIKO_CLIENT_ID env var)")
+ }
+
+ if c.IndikoURL == "" {
+ return fmt.Errorf("indiko_url is required (set INDIKO_URL env var)")
+ }
+
+ return nil
+}
+
+// findEnvFile looks for .env file in the config file's directory or current directory
+func findEnvFile(configPath string) string {
+ // If config path provided, look in its directory
+ if configPath != "" {
+ dir := filepath.Dir(configPath)
+ envPath := filepath.Join(dir, ".env")
+ if _, err := os.Stat(envPath); err == nil {
+ return envPath
+ }
+ }
+
+ // Look in current directory
+ if _, err := os.Stat(".env"); err == nil {
+ return ".env"
+ }
+
+ return ""
+}
+
+// applyEnvOverrides applies environment variable overrides to config
+func applyEnvOverrides(cfg *Config) {
+ if v := os.Getenv("ORIGIN"); v != "" {
+ cfg.Origin = v
+ }
+ if v := os.Getenv("HOST"); v != "" {
+ cfg.Host = v
+ }
+ if v := os.Getenv("PORT"); v != "" {
+ if port, err := strconv.Atoi(v); err == nil {
+ cfg.Port = port
+ }
+ }
+ if v := os.Getenv("NODE_ENV"); v != "" {
+ cfg.Env = v
+ }
+ if v := os.Getenv("LOG_LEVEL"); v != "" {
+ cfg.LogLevel = v
+ }
+ if v := os.Getenv("DATABASE_PATH"); v != "" {
+ cfg.DatabasePath = v
+ }
+ if v := os.Getenv("INDIKO_URL"); v != "" {
+ cfg.IndikoURL = v
+ }
+ if v := os.Getenv("INDIKO_CLIENT_ID"); v != "" {
+ cfg.IndikoClientID = v
+ }
+ if v := os.Getenv("INDIKO_CLIENT_SECRET"); v != "" {
+ cfg.IndikoClientSecret = v
+ }
+ if v := os.Getenv("OAUTH_CALLBACK_URL"); v != "" {
+ cfg.OAuthCallbackURL = v
+ }
+ if v := os.Getenv("SESSION_SECRET"); v != "" {
+ cfg.SessionSecret = v
+ }
+ if v := os.Getenv("SESSION_COOKIE_NAME"); v != "" {
+ cfg.SessionCookieName = v
+ }
+}
diff --git a/engine/executor.go b/engine/executor.go
new file mode 100644
index 0000000..42a0718
--- /dev/null
+++ b/engine/executor.go
@@ -0,0 +1,218 @@
+package engine
+
+import (
+ "context"
+ "encoding/json"
+ "fmt"
+ "time"
+
+ "github.com/google/uuid"
+ "github.com/kierank/pipes/nodes"
+ "github.com/kierank/pipes/store"
+)
+
+type PipeConfig struct {
+ Version string `json:"version"`
+ Nodes []Node `json:"nodes"`
+ Connections []Connection `json:"connections"`
+ Settings Settings `json:"settings"`
+}
+
+type Node struct {
+ ID string `json:"id"`
+ Type string `json:"type"`
+ Position Position `json:"position"`
+ Config map[string]interface{} `json:"config"`
+ Label string `json:"label,omitempty"`
+}
+
+type Connection struct {
+ ID string `json:"id"`
+ Source string `json:"source"`
+ Target string `json:"target"`
+ SourceHandle string `json:"sourceHandle,omitempty"`
+ TargetHandle string `json:"targetHandle,omitempty"`
+}
+
+type Position struct {
+ X float64 `json:"x"`
+ Y float64 `json:"y"`
+}
+
+type Settings struct {
+ Schedule string `json:"schedule,omitempty"`
+ Enabled bool `json:"enabled"`
+ Timeout int `json:"timeout,omitempty"`
+ RetryConfig *RetryConfig `json:"retryConfig,omitempty"`
+}
+
+type RetryConfig struct {
+ MaxRetries int `json:"maxRetries"`
+ BackoffMs int `json:"backoffMs"`
+}
+
+type Executor struct {
+ db *store.DB
+ registry *Registry
+}
+
+func NewExecutor(db *store.DB) *Executor {
+ return &Executor{
+ db: db,
+ registry: NewRegistry(),
+ }
+}
+
+func (e *Executor) Execute(ctx context.Context, pipeID string, triggerType string) (string, error) {
+ executionID := uuid.New().String()
+ startedAt := time.Now().Unix()
+
+ // Create execution record
+ if err := e.db.CreateExecution(executionID, pipeID, triggerType, startedAt); err != nil {
+ return "", fmt.Errorf("create execution: %w", err)
+ }
+
+ // Fetch pipe configuration
+ pipe, err := e.db.GetPipe(pipeID)
+ if err != nil {
+ return "", fmt.Errorf("get pipe: %w", err)
+ }
+
+ if pipe == nil {
+ return "", fmt.Errorf("pipe not found: %s", pipeID)
+ }
+
+ var config PipeConfig
+ if err := json.Unmarshal([]byte(pipe.Config), &config); err != nil {
+ return "", fmt.Errorf("parse config: %w", err)
+ }
+
+ // Execute pipeline
+ itemCount, err := e.executePipeline(ctx, executionID, pipeID, &config)
+
+ completedAt := time.Now().Unix()
+ durationMs := (completedAt - startedAt) * 1000
+
+ if err != nil {
+ e.db.UpdateExecutionFailed(executionID, completedAt, durationMs, err.Error())
+ return executionID, err
+ }
+
+ e.db.UpdateExecutionSuccess(executionID, completedAt, durationMs, itemCount)
+ return executionID, nil
+}
+
+func (e *Executor) executePipeline(ctx context.Context, executionID, pipeID string, config *PipeConfig) (int, error) {
+ // Topological sort to determine execution order
+ order, err := topologicalSort(config.Nodes, config.Connections)
+ if err != nil {
+ return 0, fmt.Errorf("topological sort: %w", err)
+ }
+
+ nodeResults := make(map[string][]interface{})
+ execCtx := nodes.NewContext(executionID, pipeID, e.db)
+
+ for _, nodeID := range order {
+ node := findNode(config.Nodes, nodeID)
+ if node == nil {
+ continue
+ }
+
+ // Get node implementation
+ nodeImpl, err := e.registry.Get(node.Type)
+ if err != nil {
+ return 0, fmt.Errorf("get node type %s: %w", node.Type, err)
+ }
+
+ // Gather inputs from connected nodes
+ inputs := e.gatherInputs(nodeID, config.Connections, nodeResults)
+
+ // Execute node
+ output, err := nodeImpl.Execute(ctx, node.Config, inputs, execCtx)
+ if err != nil {
+ e.db.LogExecution(executionID, nodeID, "error", fmt.Sprintf("Execution failed: %v", err))
+ return 0, fmt.Errorf("node %s (%s): %w", nodeID, node.Type, err)
+ }
+
+ nodeResults[nodeID] = output
+ e.db.LogExecution(executionID, nodeID, "info", fmt.Sprintf("Processed %d items", len(output)))
+ }
+
+ // Return item count from last node
+ if len(order) == 0 {
+ return 0, nil
+ }
+
+ lastNodeID := order[len(order)-1]
+ finalOutput := nodeResults[lastNodeID]
+ return len(finalOutput), nil
+}
+
+func (e *Executor) gatherInputs(nodeID string, connections []Connection, nodeResults map[string][]interface{}) [][]interface{} {
+ var inputs [][]interface{}
+
+ for _, conn := range connections {
+ if conn.Target == nodeID {
+ if result, ok := nodeResults[conn.Source]; ok {
+ inputs = append(inputs, result)
+ }
+ }
+ }
+
+ return inputs
+}
+
+func topologicalSort(nodes []Node, connections []Connection) ([]string, error) {
+ // Kahn's algorithm for topological sorting
+ inDegree := make(map[string]int)
+ adjacency := make(map[string][]string)
+
+ // Initialize
+ for _, n := range nodes {
+ inDegree[n.ID] = 0
+ adjacency[n.ID] = []string{}
+ }
+
+ // Build graph
+ for _, c := range connections {
+ adjacency[c.Source] = append(adjacency[c.Source], c.Target)
+ inDegree[c.Target]++
+ }
+
+ // Find sources (nodes with no incoming edges)
+ queue := []string{}
+ for id, degree := range inDegree {
+ if degree == 0 {
+ queue = append(queue, id)
+ }
+ }
+
+ sorted := []string{}
+ for len(queue) > 0 {
+ node := queue[0]
+ queue = queue[1:]
+ sorted = append(sorted, node)
+
+ for _, neighbor := range adjacency[node] {
+ inDegree[neighbor]--
+ if inDegree[neighbor] == 0 {
+ queue = append(queue, neighbor)
+ }
+ }
+ }
+
+ if len(sorted) != len(nodes) {
+ return nil, fmt.Errorf("pipeline contains a cycle")
+ }
+
+ return sorted, nil
+}
+
+func findNode(nodes []Node, id string) *Node {
+ for i := range nodes {
+ if nodes[i].ID == id {
+ return &nodes[i]
+ }
+ }
+ return nil
+}
diff --git a/engine/registry.go b/engine/registry.go
new file mode 100644
index 0000000..2d6c926
--- /dev/null
+++ b/engine/registry.go
@@ -0,0 +1,59 @@
+package engine
+
+import (
+ "fmt"
+ "sync"
+
+ "github.com/kierank/pipes/nodes"
+ "github.com/kierank/pipes/nodes/sources"
+ "github.com/kierank/pipes/nodes/transforms"
+)
+
+type Registry struct {
+ mu sync.RWMutex
+ nodeImpls map[string]nodes.Node
+}
+
+func NewRegistry() *Registry {
+ r := &Registry{
+ nodeImpls: make(map[string]nodes.Node),
+ }
+
+ // Register built-in nodes
+ r.Register(&sources.RSSSourceNode{})
+ r.Register(&transforms.FilterNode{})
+ r.Register(&transforms.SortNode{})
+ r.Register(&transforms.LimitNode{})
+
+ return r
+}
+
+func (r *Registry) Register(node nodes.Node) {
+ r.mu.Lock()
+ defer r.mu.Unlock()
+ r.nodeImpls[node.Type()] = node
+}
+
+func (r *Registry) Get(nodeType string) (nodes.Node, error) {
+ r.mu.RLock()
+ defer r.mu.RUnlock()
+
+ node, ok := r.nodeImpls[nodeType]
+ if !ok {
+ return nil, fmt.Errorf("unknown node type: %s", nodeType)
+ }
+
+ return node, nil
+}
+
+func (r *Registry) GetAll() []nodes.Node {
+ r.mu.RLock()
+ defer r.mu.RUnlock()
+
+ nodeList := make([]nodes.Node, 0, len(r.nodeImpls))
+ for _, node := range r.nodeImpls {
+ nodeList = append(nodeList, node)
+ }
+
+ return nodeList
+}
diff --git a/engine/scheduler.go b/engine/scheduler.go
new file mode 100644
index 0000000..a3c0f79
--- /dev/null
+++ b/engine/scheduler.go
@@ -0,0 +1,93 @@
+package engine
+
+import (
+ "context"
+ "time"
+
+ "github.com/charmbracelet/log"
+ "github.com/kierank/pipes/store"
+)
+
+type Scheduler struct {
+ db *store.DB
+ executor *Executor
+ ticker *time.Ticker
+ done chan struct{}
+ logger *log.Logger
+}
+
+func NewScheduler(db *store.DB, logger *log.Logger) *Scheduler {
+ return &Scheduler{
+ db: db,
+ executor: NewExecutor(db),
+ done: make(chan struct{}),
+ logger: logger,
+ }
+}
+
+func (s *Scheduler) Start() {
+ s.logger.Info("scheduler starting")
+
+ s.ticker = time.NewTicker(1 * time.Minute)
+
+ // Run immediately on start
+ go s.tick()
+
+ // Then run every minute
+ go func() {
+ for {
+ select {
+ case <-s.ticker.C:
+ s.tick()
+ case <-s.done:
+ return
+ }
+ }
+ }()
+}
+
+func (s *Scheduler) tick() {
+ ctx := context.Background()
+ now := time.Now().Unix()
+
+ jobs, err := s.db.GetDueJobs(now)
+ if err != nil {
+ s.logger.Error("error fetching jobs", "error", err)
+ return
+ }
+
+ if len(jobs) > 0 {
+ s.logger.Info("found jobs to execute", "count", len(jobs))
+ }
+
+ for _, job := range jobs {
+ if err := s.executeJob(ctx, job); err != nil {
+ s.logger.Error("job execution failed", "job_id", job.ID, "error", err)
+ }
+ }
+}
+
+func (s *Scheduler) executeJob(ctx context.Context, job *store.ScheduledJob) error {
+ // Execute pipeline
+ _, err := s.executor.Execute(ctx, job.PipeID, "scheduled")
+ if err != nil {
+ s.logger.Error("pipeline execution failed", "pipe_id", job.PipeID, "error", err)
+ }
+
+ // Calculate next run time (simplified: add 1 hour for now)
+ // In production, use a proper cron parser
+ nextRun := time.Now().Add(1 * time.Hour).Unix()
+
+ // Update job
+ now := time.Now().Unix()
+ return s.db.UpdateJobAfterRun(job.ID, now, nextRun)
+}
+
+func (s *Scheduler) Stop() {
+ s.logger.Info("scheduler stopping")
+ if s.ticker != nil {
+ s.ticker.Stop()
+ }
+ close(s.done)
+ s.logger.Info("scheduler stopped")
+}
diff --git a/go.mod b/go.mod
new file mode 100644
index 0000000..adaba58
--- /dev/null
+++ b/go.mod
@@ -0,0 +1,40 @@
+module github.com/kierank/pipes
+
+go 1.24
+
+require (
+ github.com/google/uuid v1.6.0
+ github.com/gorilla/sessions v1.4.0
+ github.com/mattn/go-sqlite3 v1.14.24
+ github.com/mmcdole/gofeed v1.3.0
+)
+
+require (
+ github.com/PuerkitoBio/goquery v1.8.0 // indirect
+ github.com/andybalholm/cascadia v1.3.1 // indirect
+ github.com/aymanbagabas/go-osc52/v2 v2.0.1 // indirect
+ github.com/charmbracelet/colorprofile v0.2.3-0.20250311203215-f60798e515dc // indirect
+ github.com/charmbracelet/lipgloss v1.1.0 // indirect
+ github.com/charmbracelet/log v0.4.2 // indirect
+ github.com/charmbracelet/x/ansi v0.8.0 // indirect
+ github.com/charmbracelet/x/cellbuf v0.0.13-0.20250311204145-2c3ea96c31dd // indirect
+ github.com/charmbracelet/x/term v0.2.1 // indirect
+ github.com/go-logfmt/logfmt v0.6.0 // indirect
+ github.com/gorilla/securecookie v1.1.2 // indirect
+ github.com/joho/godotenv v1.5.1 // indirect
+ github.com/json-iterator/go v1.1.12 // indirect
+ github.com/lucasb-eyer/go-colorful v1.2.0 // indirect
+ github.com/mattn/go-isatty v0.0.20 // indirect
+ github.com/mattn/go-runewidth v0.0.16 // indirect
+ github.com/mmcdole/goxpp v1.1.1-0.20240225020742-a0c311522b23 // indirect
+ github.com/modern-go/concurrent v0.0.0-20180306012644-bacd9c7ef1dd // indirect
+ github.com/modern-go/reflect2 v1.0.2 // indirect
+ github.com/muesli/termenv v0.16.0 // indirect
+ github.com/rivo/uniseg v0.4.7 // indirect
+ github.com/xo/terminfo v0.0.0-20220910002029-abceb7e1c41e // indirect
+ golang.org/x/exp v0.0.0-20231006140011-7918f672742d // indirect
+ golang.org/x/net v0.4.0 // indirect
+ golang.org/x/sys v0.30.0 // indirect
+ golang.org/x/text v0.5.0 // indirect
+ gopkg.in/yaml.v3 v3.0.1 // indirect
+)
diff --git a/go.sum b/go.sum
new file mode 100644
index 0000000..d473808
--- /dev/null
+++ b/go.sum
@@ -0,0 +1,84 @@
+github.com/PuerkitoBio/goquery v1.8.0 h1:PJTF7AmFCFKk1N6V6jmKfrNH9tV5pNE6lZMkG0gta/U=
+github.com/PuerkitoBio/goquery v1.8.0/go.mod h1:ypIiRMtY7COPGk+I/YbZLbxsxn9g5ejnI2HSMtkjZvI=
+github.com/andybalholm/cascadia v1.3.1 h1:nhxRkql1kdYCc8Snf7D5/D3spOX+dBgjA6u8x004T2c=
+github.com/andybalholm/cascadia v1.3.1/go.mod h1:R4bJ1UQfqADjvDa4P6HZHLh/3OxWWEqc0Sk8XGwHqvA=
+github.com/aymanbagabas/go-osc52/v2 v2.0.1 h1:HwpRHbFMcZLEVr42D4p7XBqjyuxQH5SMiErDT4WkJ2k=
+github.com/aymanbagabas/go-osc52/v2 v2.0.1/go.mod h1:uYgXzlJ7ZpABp8OJ+exZzJJhRNQ2ASbcXHWsFqH8hp8=
+github.com/charmbracelet/colorprofile v0.2.3-0.20250311203215-f60798e515dc h1:4pZI35227imm7yK2bGPcfpFEmuY1gc2YSTShr4iJBfs=
+github.com/charmbracelet/colorprofile v0.2.3-0.20250311203215-f60798e515dc/go.mod h1:X4/0JoqgTIPSFcRA/P6INZzIuyqdFY5rm8tb41s9okk=
+github.com/charmbracelet/lipgloss v1.1.0 h1:vYXsiLHVkK7fp74RkV7b2kq9+zDLoEU4MZoFqR/noCY=
+github.com/charmbracelet/lipgloss v1.1.0/go.mod h1:/6Q8FR2o+kj8rz4Dq0zQc3vYf7X+B0binUUBwA0aL30=
+github.com/charmbracelet/log v0.4.2 h1:hYt8Qj6a8yLnvR+h7MwsJv/XvmBJXiueUcI3cIxsyig=
+github.com/charmbracelet/log v0.4.2/go.mod h1:qifHGX/tc7eluv2R6pWIpyHDDrrb/AG71Pf2ysQu5nw=
+github.com/charmbracelet/x/ansi v0.8.0 h1:9GTq3xq9caJW8ZrBTe0LIe2fvfLR/bYXKTx2llXn7xE=
+github.com/charmbracelet/x/ansi v0.8.0/go.mod h1:wdYl/ONOLHLIVmQaxbIYEC/cRKOQyjTkowiI4blgS9Q=
+github.com/charmbracelet/x/cellbuf v0.0.13-0.20250311204145-2c3ea96c31dd h1:vy0GVL4jeHEwG5YOXDmi86oYw2yuYUGqz6a8sLwg0X8=
+github.com/charmbracelet/x/cellbuf v0.0.13-0.20250311204145-2c3ea96c31dd/go.mod h1:xe0nKWGd3eJgtqZRaN9RjMtK7xUYchjzPr7q6kcvCCs=
+github.com/charmbracelet/x/term v0.2.1 h1:AQeHeLZ1OqSXhrAWpYUtZyX1T3zVxfpZuEQMIQaGIAQ=
+github.com/charmbracelet/x/term v0.2.1/go.mod h1:oQ4enTYFV7QN4m0i9mzHrViD7TQKvNEEkHUMCmsxdUg=
+github.com/davecgh/go-spew v1.1.0/go.mod h1:J7Y8YcW2NihsgmVo/mv3lAwl/skON4iLHjSsI+c5H38=
+github.com/davecgh/go-spew v1.1.1 h1:vj9j/u1bqnvCEfJOwUhtlOARqs3+rkHYY13jYWTU97c=
+github.com/davecgh/go-spew v1.1.1/go.mod h1:J7Y8YcW2NihsgmVo/mv3lAwl/skON4iLHjSsI+c5H38=
+github.com/go-logfmt/logfmt v0.6.0 h1:wGYYu3uicYdqXVgoYbvnkrPVXkuLM1p1ifugDMEdRi4=
+github.com/go-logfmt/logfmt v0.6.0/go.mod h1:WYhtIu8zTZfxdn5+rREduYbwxfcBr/Vr6KEVveWlfTs=
+github.com/google/gofuzz v1.0.0/go.mod h1:dBl0BpW6vV/+mYPU4Po3pmUjxk6FQPldtuIdl/M65Eg=
+github.com/google/gofuzz v1.2.0 h1:xRy4A+RhZaiKjJ1bPfwQ8sedCA+YS2YcCHW6ec7JMi0=
+github.com/google/gofuzz v1.2.0/go.mod h1:dBl0BpW6vV/+mYPU4Po3pmUjxk6FQPldtuIdl/M65Eg=
+github.com/google/uuid v1.6.0 h1:NIvaJDMOsjHA8n1jAhLSgzrAzy1Hgr+hNrb57e+94F0=
+github.com/google/uuid v1.6.0/go.mod h1:TIyPZe4MgqvfeYDBFedMoGGpEw/LqOeaOT+nhxU+yHo=
+github.com/gorilla/securecookie v1.1.2 h1:YCIWL56dvtr73r6715mJs5ZvhtnY73hBvEF8kXD8ePA=
+github.com/gorilla/securecookie v1.1.2/go.mod h1:NfCASbcHqRSY+3a8tlWJwsQap2VX5pwzwo4h3eOamfo=
+github.com/gorilla/sessions v1.4.0 h1:kpIYOp/oi6MG/p5PgxApU8srsSw9tuFbt46Lt7auzqQ=
+github.com/gorilla/sessions v1.4.0/go.mod h1:FLWm50oby91+hl7p/wRxDth9bWSuk0qVL2emc7lT5ik=
+github.com/joho/godotenv v1.5.1 h1:7eLL/+HRGLY0ldzfGMeQkb7vMd0as4CfYvUVzLqw0N0=
+github.com/joho/godotenv v1.5.1/go.mod h1:f4LDr5Voq0i2e/R5DDNOoa2zzDfwtkZa6DnEwAbqwq4=
+github.com/json-iterator/go v1.1.12 h1:PV8peI4a0ysnczrg+LtxykD8LfKY9ML6u2jnxaEnrnM=
+github.com/json-iterator/go v1.1.12/go.mod h1:e30LSqwooZae/UwlEbR2852Gd8hjQvJoHmT4TnhNGBo=
+github.com/lucasb-eyer/go-colorful v1.2.0 h1:1nnpGOrhyZZuNyfu1QjKiUICQ74+3FNCN69Aj6K7nkY=
+github.com/lucasb-eyer/go-colorful v1.2.0/go.mod h1:R4dSotOR9KMtayYi1e77YzuveK+i7ruzyGqttikkLy0=
+github.com/mattn/go-isatty v0.0.20 h1:xfD0iDuEKnDkl03q4limB+vH+GxLEtL/jb4xVJSWWEY=
+github.com/mattn/go-isatty v0.0.20/go.mod h1:W+V8PltTTMOvKvAeJH7IuucS94S2C6jfK/D7dTCTo3Y=
+github.com/mattn/go-runewidth v0.0.16 h1:E5ScNMtiwvlvB5paMFdw9p4kSQzbXFikJ5SQO6TULQc=
+github.com/mattn/go-runewidth v0.0.16/go.mod h1:Jdepj2loyihRzMpdS35Xk/zdY8IAYHsh153qUoGf23w=
+github.com/mattn/go-sqlite3 v1.14.24 h1:tpSp2G2KyMnnQu99ngJ47EIkWVmliIizyZBfPrBWDRM=
+github.com/mattn/go-sqlite3 v1.14.24/go.mod h1:Uh1q+B4BYcTPb+yiD3kU8Ct7aC0hY9fxUwlHK0RXw+Y=
+github.com/mmcdole/gofeed v1.3.0 h1:5yn+HeqlcvjMeAI4gu6T+crm7d0anY85+M+v6fIFNG4=
+github.com/mmcdole/gofeed v1.3.0/go.mod h1:9TGv2LcJhdXePDzxiuMnukhV2/zb6VtnZt1mS+SjkLE=
+github.com/mmcdole/goxpp v1.1.1-0.20240225020742-a0c311522b23 h1:Zr92CAlFhy2gL+V1F+EyIuzbQNbSgP4xhTODZtrXUtk=
+github.com/mmcdole/goxpp v1.1.1-0.20240225020742-a0c311522b23/go.mod h1:v+25+lT2ViuQ7mVxcncQ8ch1URund48oH+jhjiwEgS8=
+github.com/modern-go/concurrent v0.0.0-20180228061459-e0a39a4cb421/go.mod h1:6dJC0mAP4ikYIbvyc7fijjWJddQyLn8Ig3JB5CqoB9Q=
+github.com/modern-go/concurrent v0.0.0-20180306012644-bacd9c7ef1dd h1:TRLaZ9cD/w8PVh93nsPXa1VrQ6jlwL5oN8l14QlcNfg=
+github.com/modern-go/concurrent v0.0.0-20180306012644-bacd9c7ef1dd/go.mod h1:6dJC0mAP4ikYIbvyc7fijjWJddQyLn8Ig3JB5CqoB9Q=
+github.com/modern-go/reflect2 v1.0.2 h1:xBagoLtFs94CBntxluKeaWgTMpvLxC4ur3nMaC9Gz0M=
+github.com/modern-go/reflect2 v1.0.2/go.mod h1:yWuevngMOJpCy52FWWMvUC8ws7m/LJsjYzDa0/r8luk=
+github.com/muesli/termenv v0.16.0 h1:S5AlUN9dENB57rsbnkPyfdGuWIlkmzJjbFf0Tf5FWUc=
+github.com/muesli/termenv v0.16.0/go.mod h1:ZRfOIKPFDYQoDFF4Olj7/QJbW60Ol/kL1pU3VfY/Cnk=
+github.com/pmezard/go-difflib v1.0.0 h1:4DBwDE0NGyQoBHbLQYPwSUPoCMWR5BEzIk/f1lZbAQM=
+github.com/pmezard/go-difflib v1.0.0/go.mod h1:iKH77koFhYxTK1pcRnkKkqfTogsbg7gZNVY4sRDYZ/4=
+github.com/rivo/uniseg v0.2.0/go.mod h1:J6wj4VEh+S6ZtnVlnTBMWIodfgj8LQOQFoIToxlJtxc=
+github.com/rivo/uniseg v0.4.7 h1:WUdvkW8uEhrYfLC4ZzdpI2ztxP1I582+49Oc5Mq64VQ=
+github.com/rivo/uniseg v0.4.7/go.mod h1:FN3SvrM+Zdj16jyLfmOkMNblXMcoc8DfTHruCPUcx88=
+github.com/stretchr/objx v0.1.0/go.mod h1:HFkY916IF+rwdDfMAkV7OtwuqBVzrE8GR6GFx+wExME=
+github.com/stretchr/testify v1.3.0/go.mod h1:M5WIy9Dh21IEIfnGCwXGc5bZfKNJtfHm1UVUgZn+9EI=
+github.com/stretchr/testify v1.8.1 h1:w7B6lhMri9wdJUVmEZPGGhZzrYTPvgJArz7wNPgYKsk=
+github.com/stretchr/testify v1.8.1/go.mod h1:w2LPCIKwWwSfY2zedu0+kehJoqGctiVI29o6fzry7u4=
+github.com/xo/terminfo v0.0.0-20220910002029-abceb7e1c41e h1:JVG44RsyaB9T2KIHavMF/ppJZNG9ZpyihvCd0w101no=
+github.com/xo/terminfo v0.0.0-20220910002029-abceb7e1c41e/go.mod h1:RbqR21r5mrJuqunuUZ/Dhy/avygyECGrLceyNeo4LiM=
+golang.org/x/exp v0.0.0-20231006140011-7918f672742d h1:jtJma62tbqLibJ5sFQz8bKtEM8rJBtfilJ2qTU199MI=
+golang.org/x/exp v0.0.0-20231006140011-7918f672742d/go.mod h1:ldy0pHrwJyGW56pPQzzkH36rKxoZW1tw7ZJpeKx+hdo=
+golang.org/x/net v0.0.0-20210916014120-12bc252f5db8/go.mod h1:9nx3DQGgdP8bBQD5qxJ1jj9UTztislL4KSBs9R2vV5Y=
+golang.org/x/net v0.4.0 h1:Q5QPcMlvfxFTAPV0+07Xz/MpK9NTXu2VDUuy0FeMfaU=
+golang.org/x/net v0.4.0/go.mod h1:MBQ8lrhLObU/6UmLb4fmbmk5OcyYmqtbGd/9yIeKjEE=
+golang.org/x/sys v0.0.0-20201119102817-f84b799fce68/go.mod h1:h1NjWce9XRLGQEsW7wpKNCjG9DtNlClVuFLEZdDNbEs=
+golang.org/x/sys v0.0.0-20210423082822-04245dca01da/go.mod h1:h1NjWce9XRLGQEsW7wpKNCjG9DtNlClVuFLEZdDNbEs=
+golang.org/x/sys v0.6.0/go.mod h1:oPkhp1MJrh7nUepCBck5+mAzfO9JrbApNNgaTdGDITg=
+golang.org/x/sys v0.30.0 h1:QjkSwP/36a20jFYWkSue1YwXzLmsV5Gfq7Eiy72C1uc=
+golang.org/x/sys v0.30.0/go.mod h1:/VUhepiaJMQUp4+oa/7Zr1D23ma6VTLIYjOOTFZPUcA=
+golang.org/x/term v0.0.0-20201126162022-7de9c90e9dd1/go.mod h1:bj7SfCRtBDWHUb9snDiAeCFNEtKQo2Wmx5Cou7ajbmo=
+golang.org/x/text v0.3.6/go.mod h1:5Zoc/QRtKVWzQhOtBMvqHzDpF6irO9z98xDceosuGiQ=
+golang.org/x/text v0.5.0 h1:OLmvp0KP+FVG99Ct/qFiL/Fhk4zp4QQnZ7b2U+5piUM=
+golang.org/x/text v0.5.0/go.mod h1:mrYo+phRRbMaCq/xk9113O4dZlRixOauAjOtrjsXDZ8=
+golang.org/x/tools v0.0.0-20180917221912-90fa682c2a6e/go.mod h1:n7NCudcB/nEzxVGmLbDWY5pfWTLqBcC2KZ6jyYvM4mQ=
+gopkg.in/check.v1 v0.0.0-20161208181325-20d25e280405/go.mod h1:Co6ibVJAznAaIkqp8huTwlJQCZ016jof/cbN4VW5Yz0=
+gopkg.in/yaml.v3 v3.0.1 h1:fxVm/GzAzEWqLHuvctI91KS9hhNmmWOoWu0XTYJS7CA=
+gopkg.in/yaml.v3 v3.0.1/go.mod h1:K4uyk7z7BCEPqu6E+C64Yfv1cQ7kz7rIZviUmN+EgEM=
diff --git a/main.go b/main.go
new file mode 100644
index 0000000..d03d498
--- /dev/null
+++ b/main.go
@@ -0,0 +1,243 @@
+package main
+
+import (
+ "context"
+ "crypto/rand"
+ "encoding/base64"
+ "errors"
+ "fmt"
+ "net/http"
+ "os"
+ "os/signal"
+ "strings"
+ "syscall"
+ "time"
+
+ "github.com/charmbracelet/log"
+ "github.com/kierank/pipes/config"
+ "github.com/kierank/pipes/engine"
+ "github.com/kierank/pipes/store"
+ "github.com/kierank/pipes/web"
+)
+
+var (
+ version = "dev"
+ commitHash = "dev"
+ logger *log.Logger
+)
+
+func main() {
+ // Initialize logger with default level
+ logger = log.NewWithOptions(os.Stderr, log.Options{
+ ReportTimestamp: true,
+ Level: log.InfoLevel,
+ })
+
+ if len(os.Args) < 2 {
+ printUsage()
+ os.Exit(1)
+ }
+
+ command := os.Args[1]
+
+ switch command {
+ case "serve":
+ configPath := ""
+ // Check for -c or --config flag
+ for i := 2; i < len(os.Args); i++ {
+ if (os.Args[i] == "-c" || os.Args[i] == "--config") && i+1 < len(os.Args) {
+ configPath = os.Args[i+1]
+ break
+ }
+ }
+ serve(configPath)
+ case "init":
+ initConfig()
+ case "help", "--help", "-h":
+ printUsage()
+ case "version", "--version", "-v":
+ fmt.Printf("pipes %s (%s)\n", version, commitHash)
+ default:
+ fmt.Printf("Unknown command: %s\n\n", command)
+ printUsage()
+ os.Exit(1)
+ }
+}
+
+func printUsage() {
+ fmt.Println("Pipes - Visual data pipeline builder")
+ fmt.Println()
+ fmt.Println("Usage:")
+ fmt.Println(" pipes [flags]")
+ fmt.Println()
+ fmt.Println("Commands:")
+ fmt.Println(" serve Start the server")
+ fmt.Println(" init [path] Create a sample config file (default: config.yaml)")
+ fmt.Println(" version Show version information")
+ fmt.Println(" help Show this help message")
+ fmt.Println()
+ fmt.Println("Serve Flags:")
+ fmt.Println(" -c, --config PATH Path to config file (optional, uses .env if not specified)")
+ fmt.Println()
+ fmt.Println("Examples:")
+ fmt.Println(" pipes init")
+ fmt.Println(" pipes serve -c config.yaml")
+ fmt.Println(" pipes serve # Uses .env file")
+ fmt.Println()
+}
+
+func serve(configPath string) {
+ // Load configuration
+ cfg, err := config.Load(configPath)
+ if err != nil {
+ logger.Fatal("failed to load config", "error", err)
+ }
+
+ // Set log level from config
+ level := parseLogLevel(cfg.LogLevel)
+ logger.SetLevel(level)
+
+ logger.Info("starting pipes",
+ "host", cfg.Host,
+ "port", cfg.Port,
+ "db_path", cfg.DatabasePath,
+ "log_level", cfg.LogLevel,
+ )
+
+ // Initialize database
+ db, err := store.New(cfg.DatabasePath)
+ if err != nil {
+ logger.Fatal("failed to initialize database", "error", err)
+ }
+ defer db.Close()
+
+ logger.Info("database initialized successfully")
+
+ // Initialize scheduler
+ scheduler := engine.NewScheduler(db, logger)
+ scheduler.Start()
+ defer scheduler.Stop()
+
+ logger.Info("scheduler started")
+
+ // Initialize web server
+ server := web.NewServer(cfg, db, logger)
+
+ // Start server in goroutine
+ serverErr := make(chan error, 1)
+ go func() {
+ logger.Info("starting server", "address", fmt.Sprintf("%s:%d", cfg.Host, cfg.Port))
+ if err := server.Start(); err != nil && !errors.Is(err, http.ErrServerClosed) {
+ serverErr <- err
+ }
+ }()
+
+ // Wait for interrupt signal or server error
+ sigChan := make(chan os.Signal, 1)
+ signal.Notify(sigChan, os.Interrupt, syscall.SIGTERM)
+
+ select {
+ case <-sigChan:
+ logger.Info("shutting down gracefully...")
+ case err := <-serverErr:
+ logger.Fatal("server error", "error", err)
+ }
+
+ // Graceful shutdown
+ ctx, cancel := context.WithTimeout(context.Background(), 10*time.Second)
+ defer cancel()
+
+ if err := server.Shutdown(ctx); err != nil {
+ logger.Error("server shutdown error", "error", err)
+ }
+
+ logger.Info("shutdown complete")
+}
+
+func initConfig() {
+ configPath := "config.yaml"
+ if len(os.Args) > 2 {
+ configPath = os.Args[2]
+ }
+
+ if _, err := os.Stat(configPath); err == nil {
+ fmt.Printf("Config file already exists at %s\n", configPath)
+ fmt.Println("Remove it first or specify a different path:")
+ fmt.Printf(" pipes init %s.new\n", configPath)
+ os.Exit(1)
+ }
+
+ secret, err := generateSecret()
+ if err != nil {
+ logger.Fatal("failed to generate secret", "error", err)
+ }
+
+ configContent := `# Pipes Configuration
+# See https://github.com/yourusername/pipes for documentation
+
+# Server settings
+host: localhost
+port: 3001
+origin: http://localhost:3001
+env: development
+log_level: info # debug, info, warn, error, fatal
+
+# Database
+db_path: pipes.db
+
+# OAuth (Indiko)
+# Set these environment variables or replace with actual values:
+indiko_url: ${INDIKO_URL}
+indiko_client_id: ${INDIKO_CLIENT_ID}
+indiko_client_secret: ${INDIKO_CLIENT_SECRET}
+oauth_callback_url: http://localhost:3001/auth/callback
+
+# Session
+session_secret: ` + secret + `
+session_cookie_name: pipes_session
+`
+
+ if err := os.WriteFile(configPath, []byte(configContent), 0644); err != nil {
+ logger.Fatal("failed to write config", "error", err)
+ }
+
+ fmt.Printf("✓ Config file created at %s\n\n", configPath)
+ fmt.Println("Next steps:")
+ fmt.Println(" 1. Set your Indiko OAuth environment variables:")
+ fmt.Println(" export INDIKO_URL=http://localhost:3000")
+ fmt.Println(" export INDIKO_CLIENT_ID=http://localhost:3001")
+ fmt.Println()
+ fmt.Println(" 2. Or edit the config file directly to replace ${VAR} placeholders")
+ fmt.Println()
+ fmt.Println(" 3. Start the server:")
+ fmt.Printf(" pipes serve -c %s\n", configPath)
+ fmt.Println()
+ fmt.Println(" Or use environment variables with a .env file instead:")
+ fmt.Println(" cp .env.example .env")
+ fmt.Println(" pipes serve")
+}
+
+func generateSecret() (string, error) {
+ bytes := make([]byte, 32)
+ if _, err := rand.Read(bytes); err != nil {
+ return "", err
+ }
+ return base64.URLEncoding.EncodeToString(bytes), nil
+}
+
+func parseLogLevel(levelStr string) log.Level {
+ switch strings.ToLower(levelStr) {
+ case "debug":
+ return log.DebugLevel
+ case "info":
+ return log.InfoLevel
+ case "warn", "warning":
+ return log.WarnLevel
+ case "error":
+ return log.ErrorLevel
+ case "fatal":
+ return log.FatalLevel
+ default:
+ return log.InfoLevel
+ }
+}
diff --git a/nodes/node.go b/nodes/node.go
new file mode 100644
index 0000000..c71bca8
--- /dev/null
+++ b/nodes/node.go
@@ -0,0 +1,61 @@
+package nodes
+
+import (
+ "context"
+
+ "github.com/kierank/pipes/store"
+)
+
+type Node interface {
+ Type() string
+ Label() string
+ Description() string
+ Category() string // source|transform|output
+
+ Inputs() int
+ Outputs() int
+
+ Execute(ctx context.Context, config map[string]interface{}, inputs [][]interface{}, execCtx *Context) ([]interface{}, error)
+
+ ValidateConfig(config map[string]interface{}) error
+
+ GetConfigSchema() *ConfigSchema
+}
+
+type ConfigSchema struct {
+ Fields []ConfigField `json:"fields"`
+}
+
+type ConfigField struct {
+ Name string `json:"name"`
+ Label string `json:"label"`
+ Type string `json:"type"` // text|url|number|select|textarea|checkbox
+ Required bool `json:"required,omitempty"`
+ DefaultValue interface{} `json:"defaultValue,omitempty"`
+ Options []FieldOption `json:"options,omitempty"`
+ Placeholder string `json:"placeholder,omitempty"`
+ HelpText string `json:"helpText,omitempty"`
+}
+
+type FieldOption struct {
+ Value string `json:"value"`
+ Label string `json:"label"`
+}
+
+type Context struct {
+ ExecutionID string
+ PipeID string
+ DB *store.DB
+}
+
+func NewContext(executionID, pipeID string, db *store.DB) *Context {
+ return &Context{
+ ExecutionID: executionID,
+ PipeID: pipeID,
+ DB: db,
+ }
+}
+
+func (c *Context) Log(nodeID, level, message string) {
+ c.DB.LogExecution(c.ExecutionID, nodeID, level, message)
+}
diff --git a/nodes/sources/rss.go b/nodes/sources/rss.go
new file mode 100644
index 0000000..fa9104b
--- /dev/null
+++ b/nodes/sources/rss.go
@@ -0,0 +1,93 @@
+package sources
+
+import (
+ "context"
+ "fmt"
+
+ "github.com/mmcdole/gofeed"
+
+ "github.com/kierank/pipes/nodes"
+)
+
+type RSSSourceNode struct{}
+
+func (n *RSSSourceNode) Type() string { return "rss-source" }
+func (n *RSSSourceNode) Label() string { return "RSS Feed" }
+func (n *RSSSourceNode) Description() string { return "Fetch items from an RSS or Atom feed" }
+func (n *RSSSourceNode) Category() string { return "source" }
+func (n *RSSSourceNode) Inputs() int { return 0 }
+func (n *RSSSourceNode) Outputs() int { return 1 }
+
+func (n *RSSSourceNode) Execute(ctx context.Context, config map[string]interface{}, inputs [][]interface{}, execCtx *nodes.Context) ([]interface{}, error) {
+ url, ok := config["url"].(string)
+ if !ok || url == "" {
+ return nil, fmt.Errorf("url is required")
+ }
+
+ execCtx.Log("rss-source", "info", fmt.Sprintf("Fetching %s", url))
+
+ // Parse feed
+ fp := gofeed.NewParser()
+ feed, err := fp.ParseURLWithContext(url, ctx)
+ if err != nil {
+ return nil, fmt.Errorf("parse feed: %w", err)
+ }
+
+ // Convert feed items to generic interface{} slices
+ var items []interface{}
+ for _, item := range feed.Items {
+ items = append(items, map[string]interface{}{
+ "title": item.Title,
+ "description": item.Description,
+ "link": item.Link,
+ "author": item.Author,
+ "published": item.Published,
+ "updated": item.Updated,
+ "guid": item.GUID,
+ "categories": item.Categories,
+ })
+ }
+
+ // Apply limit if specified
+ if limit, ok := config["limit"].(float64); ok && limit > 0 {
+ if int(limit) < len(items) {
+ items = items[:int(limit)]
+ }
+ }
+
+ execCtx.Log("rss-source", "info", fmt.Sprintf("Retrieved %d items", len(items)))
+
+ return items, nil
+}
+
+func (n *RSSSourceNode) ValidateConfig(config map[string]interface{}) error {
+ url, ok := config["url"].(string)
+ if !ok || url == "" {
+ return fmt.Errorf("url is required")
+ }
+
+ return nil
+}
+
+func (n *RSSSourceNode) GetConfigSchema() *nodes.ConfigSchema {
+ return &nodes.ConfigSchema{
+ Fields: []nodes.ConfigField{
+ {
+ Name: "url",
+ Label: "Feed URL",
+ Type: "url",
+ Required: true,
+ Placeholder: "https://example.com/feed.xml",
+ HelpText: "URL of the RSS or Atom feed",
+ },
+ {
+ Name: "limit",
+ Label: "Item Limit",
+ Type: "number",
+ Required: false,
+ DefaultValue: 50,
+ HelpText: "Maximum number of items to fetch",
+ },
+ },
+ }
+}
diff --git a/nodes/transforms/filter.go b/nodes/transforms/filter.go
new file mode 100644
index 0000000..559915f
--- /dev/null
+++ b/nodes/transforms/filter.go
@@ -0,0 +1,123 @@
+package transforms
+
+import (
+ "context"
+ "fmt"
+ "regexp"
+ "strings"
+
+ "github.com/kierank/pipes/nodes"
+)
+
+type FilterNode struct{}
+
+func (n *FilterNode) Type() string { return "filter" }
+func (n *FilterNode) Label() string { return "Filter" }
+func (n *FilterNode) Description() string { return "Filter items based on conditions" }
+func (n *FilterNode) Category() string { return "transform" }
+func (n *FilterNode) Inputs() int { return 1 }
+func (n *FilterNode) Outputs() int { return 1 }
+
+func (n *FilterNode) Execute(ctx context.Context, config map[string]interface{}, inputs [][]interface{}, execCtx *nodes.Context) ([]interface{}, error) {
+ if len(inputs) == 0 {
+ return []interface{}{}, nil
+ }
+
+ items := inputs[0]
+
+ field, _ := config["field"].(string)
+ operator, _ := config["operator"].(string)
+ value, _ := config["value"].(string)
+
+ if field == "" || operator == "" {
+ return items, nil
+ }
+
+ var filtered []interface{}
+ for _, item := range items {
+ if matchesFilter(item, field, operator, value) {
+ filtered = append(filtered, item)
+ }
+ }
+
+ execCtx.Log("filter", "info", fmt.Sprintf("Filtered %d -> %d items", len(items), len(filtered)))
+
+ return filtered, nil
+}
+
+func matchesFilter(item interface{}, field, operator, value string) bool {
+ itemMap, ok := item.(map[string]interface{})
+ if !ok {
+ return false
+ }
+
+ fieldValue := getNestedValue(itemMap, field)
+ fieldStr := fmt.Sprintf("%v", fieldValue)
+
+ switch operator {
+ case "contains":
+ return strings.Contains(strings.ToLower(fieldStr), strings.ToLower(value))
+ case "equals":
+ return fieldStr == value
+ case "not-equals":
+ return fieldStr != value
+ case "regex":
+ matched, _ := regexp.MatchString(value, fieldStr)
+ return matched
+ default:
+ return true
+ }
+}
+
+func getNestedValue(obj map[string]interface{}, path string) interface{} {
+ parts := strings.Split(path, ".")
+ var current interface{} = obj
+
+ for _, part := range parts {
+ if m, ok := current.(map[string]interface{}); ok {
+ current = m[part]
+ } else {
+ return nil
+ }
+ }
+
+ return current
+}
+
+func (n *FilterNode) ValidateConfig(config map[string]interface{}) error {
+ return nil
+}
+
+func (n *FilterNode) GetConfigSchema() *nodes.ConfigSchema {
+ return &nodes.ConfigSchema{
+ Fields: []nodes.ConfigField{
+ {
+ Name: "field",
+ Label: "Field Path",
+ Type: "text",
+ Required: true,
+ Placeholder: "title",
+ HelpText: "Field to filter on (use dot notation for nested: author.name)",
+ },
+ {
+ Name: "operator",
+ Label: "Operator",
+ Type: "select",
+ Required: true,
+ Options: []nodes.FieldOption{
+ {Value: "contains", Label: "Contains"},
+ {Value: "equals", Label: "Equals"},
+ {Value: "not-equals", Label: "Not Equals"},
+ {Value: "regex", Label: "Regex Match"},
+ },
+ },
+ {
+ Name: "value",
+ Label: "Value",
+ Type: "text",
+ Required: true,
+ Placeholder: "search term",
+ },
+ },
+ }
+}
diff --git a/nodes/transforms/limit.go b/nodes/transforms/limit.go
new file mode 100644
index 0000000..54f22a2
--- /dev/null
+++ b/nodes/transforms/limit.go
@@ -0,0 +1,54 @@
+package transforms
+
+import (
+ "context"
+ "fmt"
+
+ "github.com/kierank/pipes/nodes"
+)
+
+type LimitNode struct{}
+
+func (n *LimitNode) Type() string { return "limit" }
+func (n *LimitNode) Label() string { return "Limit" }
+func (n *LimitNode) Description() string { return "Limit the number of items" }
+func (n *LimitNode) Category() string { return "transform" }
+func (n *LimitNode) Inputs() int { return 1 }
+func (n *LimitNode) Outputs() int { return 1 }
+
+func (n *LimitNode) Execute(ctx context.Context, config map[string]interface{}, inputs [][]interface{}, execCtx *nodes.Context) ([]interface{}, error) {
+ if len(inputs) == 0 {
+ return []interface{}{}, nil
+ }
+
+ items := inputs[0]
+ count, _ := config["count"].(float64)
+
+ if count <= 0 || int(count) >= len(items) {
+ return items, nil
+ }
+
+ limited := items[:int(count)]
+ execCtx.Log("limit", "info", fmt.Sprintf("Limited %d -> %d items", len(items), len(limited)))
+
+ return limited, nil
+}
+
+func (n *LimitNode) ValidateConfig(config map[string]interface{}) error {
+ return nil
+}
+
+func (n *LimitNode) GetConfigSchema() *nodes.ConfigSchema {
+ return &nodes.ConfigSchema{
+ Fields: []nodes.ConfigField{
+ {
+ Name: "count",
+ Label: "Count",
+ Type: "number",
+ Required: true,
+ DefaultValue: 10,
+ HelpText: "Maximum number of items to output",
+ },
+ },
+ }
+}
diff --git a/nodes/transforms/sort.go b/nodes/transforms/sort.go
new file mode 100644
index 0000000..7960c15
--- /dev/null
+++ b/nodes/transforms/sort.go
@@ -0,0 +1,91 @@
+package transforms
+
+import (
+ "context"
+ "fmt"
+ "sort"
+
+ "github.com/kierank/pipes/nodes"
+)
+
+type SortNode struct{}
+
+func (n *SortNode) Type() string { return "sort" }
+func (n *SortNode) Label() string { return "Sort" }
+func (n *SortNode) Description() string { return "Sort items by a field" }
+func (n *SortNode) Category() string { return "transform" }
+func (n *SortNode) Inputs() int { return 1 }
+func (n *SortNode) Outputs() int { return 1 }
+
+func (n *SortNode) Execute(ctx context.Context, config map[string]interface{}, inputs [][]interface{}, execCtx *nodes.Context) ([]interface{}, error) {
+ if len(inputs) == 0 {
+ return []interface{}{}, nil
+ }
+
+ items := inputs[0]
+ field, _ := config["field"].(string)
+ order, _ := config["order"].(string)
+
+ if field == "" {
+ return items, nil
+ }
+
+ if order == "" {
+ order = "asc"
+ }
+
+ // Create a sortable slice
+ sorted := make([]interface{}, len(items))
+ copy(sorted, items)
+
+ sort.SliceStable(sorted, func(i, j int) bool {
+ iMap, iOk := sorted[i].(map[string]interface{})
+ jMap, jOk := sorted[j].(map[string]interface{})
+
+ if !iOk || !jOk {
+ return false
+ }
+
+ iVal := fmt.Sprintf("%v", getNestedValue(iMap, field))
+ jVal := fmt.Sprintf("%v", getNestedValue(jMap, field))
+
+ if order == "desc" {
+ return iVal > jVal
+ }
+ return iVal < jVal
+ })
+
+ execCtx.Log("sort", "info", fmt.Sprintf("Sorted %d items by %s (%s)", len(sorted), field, order))
+
+ return sorted, nil
+}
+
+func (n *SortNode) ValidateConfig(config map[string]interface{}) error {
+ return nil
+}
+
+func (n *SortNode) GetConfigSchema() *nodes.ConfigSchema {
+ return &nodes.ConfigSchema{
+ Fields: []nodes.ConfigField{
+ {
+ Name: "field",
+ Label: "Field Path",
+ Type: "text",
+ Required: true,
+ Placeholder: "published",
+ HelpText: "Field to sort by",
+ },
+ {
+ Name: "order",
+ Label: "Order",
+ Type: "select",
+ Required: false,
+ DefaultValue: "asc",
+ Options: []nodes.FieldOption{
+ {Value: "asc", Label: "Ascending"},
+ {Value: "desc", Label: "Descending"},
+ },
+ },
+ },
+ }
+}
diff --git a/package.json b/package.json
deleted file mode 100644
index fa37aeb..0000000
--- a/package.json
+++ /dev/null
@@ -1,23 +0,0 @@
-{
- "name": "pipes",
- "module": "index.ts",
- "type": "module",
- "private": true,
- "scripts": {
- "dev": "bun run --hot src/index.ts",
- "start": "bun run src/index.ts",
- "format": "bun run --bun biome check --write ."
- },
- "devDependencies": {
- "@types/bun": "latest"
- },
- "peerDependencies": {
- "typescript": "^5"
- },
- "dependencies": {
- "bun-sqlite-migrations": "^1.0.2",
- "kysely": "^0.28.9",
- "kysely-bun-sqlite": "^0.4.0",
- "nanoid": "^5.1.6"
- }
-}
diff --git a/src/index.ts b/src/index.ts
deleted file mode 100644
index 1097b67..0000000
--- a/src/index.ts
+++ /dev/null
@@ -1,54 +0,0 @@
-import { env } from "bun";
-import indexHTML from "./pages/index.html";
-
-(() => {
- const required = ["ORIGIN"];
-
- const missing = required.filter((key) => !process.env[key]);
-
- if (missing.length > 0) {
- console.warn(
- `[Startup] Missing required environment variables: ${missing.join(", ")}`,
- );
- process.exit(1);
- }
-
- // Validate ORIGIN is HTTPS in production
- const origin = process.env.ORIGIN as string;
- const nodeEnv = process.env.NODE_ENV || "development";
-
- if (nodeEnv === "production" && !origin.startsWith("https://")) {
- console.error(
- `[Startup] ORIGIN must use HTTPS in production (got: ${origin})`,
- );
- process.exit(1);
- }
-
- console.log(`[Startup] Environment validated (${nodeEnv} mode)`);
-})();
-
-const server = Bun.serve({
- port: env.PORT ? Number.parseInt(env.PORT, 10) : 3000,
- routes: {
- "/": indexHTML,
- },
- development: process.env.NODE_ENV !== "production",
-});
-
-console.log(`Pipes running on ${env.ORIGIN}`)
-
-let is_shutting_down = false;
-function shutdown(sig: string) {
- if (is_shutting_down) return;
- is_shutting_down = true;
-
- console.log(`[Shutdown] triggering shutdown due to ${sig}`);
-
- server.stop();
- console.log("[Shutdown] stopped server");
-
- process.exit(0);
-}
-
-process.on("SIGTERM", () => shutdown("SIGTERM"));
-process.on("SIGINT", () => shutdown("SIGINT"));
diff --git a/src/pages/index.html b/src/pages/index.html
deleted file mode 100644
index ef08efc..0000000
--- a/src/pages/index.html
+++ /dev/null
@@ -1,15 +0,0 @@
-
-
-
-
-
-
- Pipes
-
-
-
-
- Pipes
-
-
-
\ No newline at end of file
diff --git a/src/types/env.d.ts b/src/types/env.d.ts
deleted file mode 100644
index 6caf8c6..0000000
--- a/src/types/env.d.ts
+++ /dev/null
@@ -1,7 +0,0 @@
-declare module "bun" {
- interface Env {
- ORIGIN: string;
- NODE_ENV?: "dev" | "production";
- PORT?: string;
- }
-}
diff --git a/store/db.go b/store/db.go
new file mode 100644
index 0000000..62ba1b7
--- /dev/null
+++ b/store/db.go
@@ -0,0 +1,148 @@
+package store
+
+import (
+ "database/sql"
+ "fmt"
+
+ _ "github.com/mattn/go-sqlite3"
+)
+
+type DB struct {
+ *sql.DB
+}
+
+func New(path string) (*DB, error) {
+ db, err := sql.Open("sqlite3", path)
+ if err != nil {
+ return nil, fmt.Errorf("open database: %w", err)
+ }
+
+ // Enable foreign keys and WAL mode
+ if _, err := db.Exec("PRAGMA foreign_keys = ON"); err != nil {
+ return nil, fmt.Errorf("enable foreign keys: %w", err)
+ }
+
+ if _, err := db.Exec("PRAGMA journal_mode = WAL"); err != nil {
+ return nil, fmt.Errorf("enable WAL mode: %w", err)
+ }
+
+ store := &DB{DB: db}
+
+ // Initialize schema
+ if err := store.initSchema(); err != nil {
+ return nil, fmt.Errorf("init schema: %w", err)
+ }
+
+ return store, nil
+}
+
+func (db *DB) initSchema() error {
+ schema := `
+ -- Users (OAuth profiles)
+ CREATE TABLE IF NOT EXISTS users (
+ id TEXT PRIMARY KEY,
+ indiko_sub TEXT UNIQUE NOT NULL,
+ username TEXT,
+ name TEXT,
+ email TEXT,
+ photo TEXT,
+ url TEXT,
+ role TEXT NOT NULL DEFAULT 'user',
+ created_at INTEGER NOT NULL,
+ updated_at INTEGER NOT NULL
+ );
+
+ -- Sessions (OAuth sessions)
+ CREATE TABLE IF NOT EXISTS sessions (
+ id TEXT PRIMARY KEY,
+ user_id TEXT NOT NULL REFERENCES users(id) ON DELETE CASCADE,
+ access_token TEXT NOT NULL,
+ refresh_token TEXT,
+ expires_at INTEGER NOT NULL,
+ created_at INTEGER NOT NULL
+ );
+
+ CREATE INDEX IF NOT EXISTS idx_sessions_user_id ON sessions(user_id);
+
+ -- Pipes (pipeline configurations)
+ CREATE TABLE IF NOT EXISTS pipes (
+ id TEXT PRIMARY KEY,
+ user_id TEXT NOT NULL REFERENCES users(id) ON DELETE CASCADE,
+ name TEXT NOT NULL,
+ description TEXT,
+ config TEXT NOT NULL,
+ is_public INTEGER NOT NULL DEFAULT 0,
+ created_at INTEGER NOT NULL,
+ updated_at INTEGER NOT NULL
+ );
+
+ CREATE INDEX IF NOT EXISTS idx_pipes_user_id ON pipes(user_id);
+
+ -- Scheduled jobs
+ CREATE TABLE IF NOT EXISTS scheduled_jobs (
+ id TEXT PRIMARY KEY,
+ pipe_id TEXT NOT NULL UNIQUE REFERENCES pipes(id) ON DELETE CASCADE,
+ cron_expression TEXT NOT NULL,
+ next_run_at INTEGER NOT NULL,
+ last_run_at INTEGER,
+ enabled INTEGER NOT NULL DEFAULT 1,
+ created_at INTEGER NOT NULL,
+ updated_at INTEGER NOT NULL
+ );
+
+ CREATE INDEX IF NOT EXISTS idx_jobs_next_run ON scheduled_jobs(next_run_at, enabled);
+
+ -- Execution history
+ CREATE TABLE IF NOT EXISTS pipe_executions (
+ id TEXT PRIMARY KEY,
+ pipe_id TEXT NOT NULL REFERENCES pipes(id) ON DELETE CASCADE,
+ status TEXT NOT NULL,
+ trigger_type TEXT NOT NULL,
+ started_at INTEGER NOT NULL,
+ completed_at INTEGER,
+ duration_ms INTEGER,
+ items_processed INTEGER,
+ error_message TEXT,
+ metadata TEXT
+ );
+
+ CREATE INDEX IF NOT EXISTS idx_executions_pipe_id ON pipe_executions(pipe_id);
+ CREATE INDEX IF NOT EXISTS idx_executions_status ON pipe_executions(status);
+
+ -- Execution logs (detailed step logs)
+ CREATE TABLE IF NOT EXISTS execution_logs (
+ id TEXT PRIMARY KEY,
+ execution_id TEXT NOT NULL REFERENCES pipe_executions(id) ON DELETE CASCADE,
+ node_id TEXT NOT NULL,
+ level TEXT NOT NULL,
+ message TEXT NOT NULL,
+ timestamp INTEGER NOT NULL,
+ metadata TEXT
+ );
+
+ CREATE INDEX IF NOT EXISTS idx_logs_execution_id ON execution_logs(execution_id);
+
+ -- Source cache (avoid redundant fetches)
+ CREATE TABLE IF NOT EXISTS source_cache (
+ id TEXT PRIMARY KEY,
+ pipe_id TEXT NOT NULL REFERENCES pipes(id) ON DELETE CASCADE,
+ node_id TEXT NOT NULL,
+ cache_key TEXT NOT NULL,
+ data TEXT NOT NULL,
+ etag TEXT,
+ last_modified TEXT,
+ expires_at INTEGER NOT NULL,
+ created_at INTEGER NOT NULL
+ );
+
+ CREATE INDEX IF NOT EXISTS idx_cache_pipe_node ON source_cache(pipe_id, node_id);
+ CREATE INDEX IF NOT EXISTS idx_cache_expires ON source_cache(expires_at);
+ `
+
+ _, err := db.Exec(schema)
+ if err != nil {
+ return fmt.Errorf("execute schema: %w", err)
+ }
+
+ return nil
+}
diff --git a/store/executions.go b/store/executions.go
new file mode 100644
index 0000000..b71138e
--- /dev/null
+++ b/store/executions.go
@@ -0,0 +1,221 @@
+package store
+
+import (
+ "database/sql"
+ "fmt"
+ "time"
+
+ "github.com/google/uuid"
+)
+
+type PipeExecution struct {
+ ID string
+ PipeID string
+ Status string
+ TriggerType string
+ StartedAt int64
+ CompletedAt *int64
+ DurationMs *int64
+ ItemsProcessed *int
+ ErrorMessage *string
+ Metadata *string
+}
+
+type ExecutionLog struct {
+ ID string
+ ExecutionID string
+ NodeID string
+ Level string
+ Message string
+ Timestamp int64
+ Metadata *string
+}
+
+func (db *DB) CreateExecution(id, pipeID, triggerType string, startedAt int64) error {
+ _, err := db.Exec(`
+ INSERT INTO pipe_executions (id, pipe_id, status, trigger_type, started_at)
+ VALUES (?, ?, ?, ?, ?)
+ `, id, pipeID, "running", triggerType, startedAt)
+
+ if err != nil {
+ return fmt.Errorf("insert execution: %w", err)
+ }
+
+ return nil
+}
+
+func (db *DB) UpdateExecutionSuccess(id string, completedAt, durationMs int64, itemsProcessed int) error {
+ _, err := db.Exec(`
+ UPDATE pipe_executions
+ SET status = ?, completed_at = ?, duration_ms = ?, items_processed = ?
+ WHERE id = ?
+ `, "success", completedAt, durationMs, itemsProcessed, id)
+
+ if err != nil {
+ return fmt.Errorf("update execution: %w", err)
+ }
+
+ return nil
+}
+
+func (db *DB) UpdateExecutionFailed(id string, completedAt, durationMs int64, errorMessage string) error {
+ _, err := db.Exec(`
+ UPDATE pipe_executions
+ SET status = ?, completed_at = ?, duration_ms = ?, error_message = ?
+ WHERE id = ?
+ `, "failed", completedAt, durationMs, errorMessage, id)
+
+ if err != nil {
+ return fmt.Errorf("update execution: %w", err)
+ }
+
+ return nil
+}
+
+func (db *DB) GetExecution(id string) (*PipeExecution, error) {
+ exec := &PipeExecution{}
+ var completedAt, durationMs sql.NullInt64
+ var itemsProcessed sql.NullInt64
+ var errorMessage, metadata sql.NullString
+
+ err := db.QueryRow(`
+ SELECT id, pipe_id, status, trigger_type, started_at, completed_at, duration_ms, items_processed, error_message, metadata
+ FROM pipe_executions
+ WHERE id = ?
+ `, id).Scan(&exec.ID, &exec.PipeID, &exec.Status, &exec.TriggerType, &exec.StartedAt, &completedAt, &durationMs, &itemsProcessed, &errorMessage, &metadata)
+
+ if err == sql.ErrNoRows {
+ return nil, nil
+ }
+
+ if err != nil {
+ return nil, fmt.Errorf("query execution: %w", err)
+ }
+
+ if completedAt.Valid {
+ val := completedAt.Int64
+ exec.CompletedAt = &val
+ }
+
+ if durationMs.Valid {
+ val := durationMs.Int64
+ exec.DurationMs = &val
+ }
+
+ if itemsProcessed.Valid {
+ val := int(itemsProcessed.Int64)
+ exec.ItemsProcessed = &val
+ }
+
+ if errorMessage.Valid {
+ exec.ErrorMessage = &errorMessage.String
+ }
+
+ if metadata.Valid {
+ exec.Metadata = &metadata.String
+ }
+
+ return exec, nil
+}
+
+func (db *DB) GetPipeExecutions(pipeID string, limit int) ([]*PipeExecution, error) {
+ rows, err := db.Query(`
+ SELECT id, pipe_id, status, trigger_type, started_at, completed_at, duration_ms, items_processed, error_message, metadata
+ FROM pipe_executions
+ WHERE pipe_id = ?
+ ORDER BY started_at DESC
+ LIMIT ?
+ `, pipeID, limit)
+
+ if err != nil {
+ return nil, fmt.Errorf("query executions: %w", err)
+ }
+ defer rows.Close()
+
+ var executions []*PipeExecution
+ for rows.Next() {
+ exec := &PipeExecution{}
+ var completedAt, durationMs sql.NullInt64
+ var itemsProcessed sql.NullInt64
+ var errorMessage, metadata sql.NullString
+
+ if err := rows.Scan(&exec.ID, &exec.PipeID, &exec.Status, &exec.TriggerType, &exec.StartedAt, &completedAt, &durationMs, &itemsProcessed, &errorMessage, &metadata); err != nil {
+ return nil, fmt.Errorf("scan execution: %w", err)
+ }
+
+ if completedAt.Valid {
+ val := completedAt.Int64
+ exec.CompletedAt = &val
+ }
+
+ if durationMs.Valid {
+ val := durationMs.Int64
+ exec.DurationMs = &val
+ }
+
+ if itemsProcessed.Valid {
+ val := int(itemsProcessed.Int64)
+ exec.ItemsProcessed = &val
+ }
+
+ if errorMessage.Valid {
+ exec.ErrorMessage = &errorMessage.String
+ }
+
+ if metadata.Valid {
+ exec.Metadata = &metadata.String
+ }
+
+ executions = append(executions, exec)
+ }
+
+ return executions, nil
+}
+
+func (db *DB) LogExecution(executionID, nodeID, level, message string) error {
+ logID := uuid.New().String()
+ timestamp := time.Now().Unix()
+
+ _, err := db.Exec(`
+ INSERT INTO execution_logs (id, execution_id, node_id, level, message, timestamp)
+ VALUES (?, ?, ?, ?, ?, ?)
+ `, logID, executionID, nodeID, level, message, timestamp)
+
+ if err != nil {
+ return fmt.Errorf("insert log: %w", err)
+ }
+
+ return nil
+}
+
+func (db *DB) GetExecutionLogs(executionID string) ([]*ExecutionLog, error) {
+ rows, err := db.Query(`
+ SELECT id, execution_id, node_id, level, message, timestamp, metadata
+ FROM execution_logs
+ WHERE execution_id = ?
+ ORDER BY timestamp ASC
+ `, executionID)
+
+ if err != nil {
+ return nil, fmt.Errorf("query logs: %w", err)
+ }
+ defer rows.Close()
+
+ var logs []*ExecutionLog
+ for rows.Next() {
+ log := &ExecutionLog{}
+ var metadata sql.NullString
+
+ if err := rows.Scan(&log.ID, &log.ExecutionID, &log.NodeID, &log.Level, &log.Message, &log.Timestamp, &metadata); err != nil {
+ return nil, fmt.Errorf("scan log: %w", err)
+ }
+
+ if metadata.Valid {
+ log.Metadata = &metadata.String
+ }
+
+ logs = append(logs, log)
+ }
+
+ return logs, nil
+}
diff --git a/store/pipes.go b/store/pipes.go
new file mode 100644
index 0000000..da006d6
--- /dev/null
+++ b/store/pipes.go
@@ -0,0 +1,212 @@
+package store
+
+import (
+ "database/sql"
+ "fmt"
+ "time"
+
+ "github.com/google/uuid"
+)
+
+type Pipe struct {
+ ID string
+ UserID string
+ Name string
+ Description string
+ Config string
+ IsPublic bool
+ CreatedAt int64
+ UpdatedAt int64
+}
+
+type ScheduledJob struct {
+ ID string
+ PipeID string
+ CronExpression string
+ NextRunAt int64
+ LastRunAt *int64
+ Enabled bool
+ CreatedAt int64
+ UpdatedAt int64
+}
+
+func (db *DB) CreatePipe(userID, name, description, config string, isPublic bool) (*Pipe, error) {
+ now := time.Now().Unix()
+ pipe := &Pipe{
+ ID: uuid.New().String(),
+ UserID: userID,
+ Name: name,
+ Description: description,
+ Config: config,
+ IsPublic: isPublic,
+ CreatedAt: now,
+ UpdatedAt: now,
+ }
+
+ _, err := db.Exec(`
+ INSERT INTO pipes (id, user_id, name, description, config, is_public, created_at, updated_at)
+ VALUES (?, ?, ?, ?, ?, ?, ?, ?)
+ `, pipe.ID, pipe.UserID, pipe.Name, pipe.Description, pipe.Config, btoi(pipe.IsPublic), pipe.CreatedAt, pipe.UpdatedAt)
+
+ if err != nil {
+ return nil, fmt.Errorf("insert pipe: %w", err)
+ }
+
+ return pipe, nil
+}
+
+func (db *DB) GetPipe(id string) (*Pipe, error) {
+ pipe := &Pipe{}
+ var isPublic int
+
+ err := db.QueryRow(`
+ SELECT id, user_id, name, description, config, is_public, created_at, updated_at
+ FROM pipes
+ WHERE id = ?
+ `, id).Scan(&pipe.ID, &pipe.UserID, &pipe.Name, &pipe.Description, &pipe.Config, &isPublic, &pipe.CreatedAt, &pipe.UpdatedAt)
+
+ if err == sql.ErrNoRows {
+ return nil, nil
+ }
+
+ if err != nil {
+ return nil, fmt.Errorf("query pipe: %w", err)
+ }
+
+ pipe.IsPublic = isPublic == 1
+ return pipe, nil
+}
+
+func (db *DB) GetUserPipes(userID string) ([]*Pipe, error) {
+ rows, err := db.Query(`
+ SELECT id, user_id, name, description, config, is_public, created_at, updated_at
+ FROM pipes
+ WHERE user_id = ?
+ ORDER BY updated_at DESC
+ `, userID)
+
+ if err != nil {
+ return nil, fmt.Errorf("query pipes: %w", err)
+ }
+ defer rows.Close()
+
+ var pipes []*Pipe
+ for rows.Next() {
+ pipe := &Pipe{}
+ var isPublic int
+
+ if err := rows.Scan(&pipe.ID, &pipe.UserID, &pipe.Name, &pipe.Description, &pipe.Config, &isPublic, &pipe.CreatedAt, &pipe.UpdatedAt); err != nil {
+ return nil, fmt.Errorf("scan pipe: %w", err)
+ }
+
+ pipe.IsPublic = isPublic == 1
+ pipes = append(pipes, pipe)
+ }
+
+ return pipes, nil
+}
+
+func (db *DB) UpdatePipe(pipe *Pipe) error {
+ pipe.UpdatedAt = time.Now().Unix()
+
+ _, err := db.Exec(`
+ UPDATE pipes
+ SET name = ?, description = ?, config = ?, is_public = ?, updated_at = ?
+ WHERE id = ?
+ `, pipe.Name, pipe.Description, pipe.Config, btoi(pipe.IsPublic), pipe.UpdatedAt, pipe.ID)
+
+ if err != nil {
+ return fmt.Errorf("update pipe: %w", err)
+ }
+
+ return nil
+}
+
+func (db *DB) DeletePipe(id string) error {
+ _, err := db.Exec("DELETE FROM pipes WHERE id = ?", id)
+ if err != nil {
+ return fmt.Errorf("delete pipe: %w", err)
+ }
+ return nil
+}
+
+func (db *DB) CreateScheduledJob(pipeID, cronExpression string, nextRunAt int64) (*ScheduledJob, error) {
+ now := time.Now().Unix()
+ job := &ScheduledJob{
+ ID: uuid.New().String(),
+ PipeID: pipeID,
+ CronExpression: cronExpression,
+ NextRunAt: nextRunAt,
+ Enabled: true,
+ CreatedAt: now,
+ UpdatedAt: now,
+ }
+
+ _, err := db.Exec(`
+ INSERT INTO scheduled_jobs (id, pipe_id, cron_expression, next_run_at, enabled, created_at, updated_at)
+ VALUES (?, ?, ?, ?, ?, ?, ?)
+ `, job.ID, job.PipeID, job.CronExpression, job.NextRunAt, btoi(job.Enabled), job.CreatedAt, job.UpdatedAt)
+
+ if err != nil {
+ return nil, fmt.Errorf("insert scheduled job: %w", err)
+ }
+
+ return job, nil
+}
+
+func (db *DB) GetDueJobs(now int64) ([]*ScheduledJob, error) {
+ rows, err := db.Query(`
+ SELECT id, pipe_id, cron_expression, next_run_at, last_run_at, enabled, created_at, updated_at
+ FROM scheduled_jobs
+ WHERE enabled = 1 AND next_run_at <= ?
+ `, now)
+
+ if err != nil {
+ return nil, fmt.Errorf("query due jobs: %w", err)
+ }
+ defer rows.Close()
+
+ var jobs []*ScheduledJob
+ for rows.Next() {
+ job := &ScheduledJob{}
+ var enabled int
+ var lastRunAt sql.NullInt64
+
+ if err := rows.Scan(&job.ID, &job.PipeID, &job.CronExpression, &job.NextRunAt, &lastRunAt, &enabled, &job.CreatedAt, &job.UpdatedAt); err != nil {
+ return nil, fmt.Errorf("scan job: %w", err)
+ }
+
+ job.Enabled = enabled == 1
+ if lastRunAt.Valid {
+ val := lastRunAt.Int64
+ job.LastRunAt = &val
+ }
+
+ jobs = append(jobs, job)
+ }
+
+ return jobs, nil
+}
+
+func (db *DB) UpdateJobAfterRun(id string, lastRunAt, nextRunAt int64) error {
+ now := time.Now().Unix()
+
+ _, err := db.Exec(`
+ UPDATE scheduled_jobs
+ SET last_run_at = ?, next_run_at = ?, updated_at = ?
+ WHERE id = ?
+ `, lastRunAt, nextRunAt, now, id)
+
+ if err != nil {
+ return fmt.Errorf("update job: %w", err)
+ }
+
+ return nil
+}
+
+func btoi(b bool) int {
+ if b {
+ return 1
+ }
+ return 0
+}
diff --git a/store/users.go b/store/users.go
new file mode 100644
index 0000000..8b18649
--- /dev/null
+++ b/store/users.go
@@ -0,0 +1,170 @@
+package store
+
+import (
+ "database/sql"
+ "fmt"
+ "time"
+
+ "github.com/google/uuid"
+)
+
+type User struct {
+ ID string
+ IndikoSub string
+ Username string
+ Name string
+ Email string
+ Photo string
+ URL string
+ Role string
+ CreatedAt int64
+ UpdatedAt int64
+}
+
+type Session struct {
+ ID string
+ UserID string
+ AccessToken string
+ RefreshToken string
+ ExpiresAt int64
+ CreatedAt int64
+}
+
+func (db *DB) CreateUser(indikoSub, username, name, email, photo, url string) (*User, error) {
+ now := time.Now().Unix()
+ user := &User{
+ ID: uuid.New().String(),
+ IndikoSub: indikoSub,
+ Username: username,
+ Name: name,
+ Email: email,
+ Photo: photo,
+ URL: url,
+ Role: "user",
+ CreatedAt: now,
+ UpdatedAt: now,
+ }
+
+ _, err := db.Exec(`
+ INSERT INTO users (id, indiko_sub, username, name, email, photo, url, role, created_at, updated_at)
+ VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?)
+ `, user.ID, user.IndikoSub, user.Username, user.Name, user.Email, user.Photo, user.URL, user.Role, user.CreatedAt, user.UpdatedAt)
+
+ if err != nil {
+ return nil, fmt.Errorf("insert user: %w", err)
+ }
+
+ return user, nil
+}
+
+func (db *DB) GetUserByIndikoSub(indikoSub string) (*User, error) {
+ user := &User{}
+ err := db.QueryRow(`
+ SELECT id, indiko_sub, username, name, email, photo, url, role, created_at, updated_at
+ FROM users
+ WHERE indiko_sub = ?
+ `, indikoSub).Scan(&user.ID, &user.IndikoSub, &user.Username, &user.Name, &user.Email, &user.Photo, &user.URL, &user.Role, &user.CreatedAt, &user.UpdatedAt)
+
+ if err == sql.ErrNoRows {
+ return nil, nil
+ }
+
+ if err != nil {
+ return nil, fmt.Errorf("query user: %w", err)
+ }
+
+ return user, nil
+}
+
+func (db *DB) GetUserByID(id string) (*User, error) {
+ user := &User{}
+ err := db.QueryRow(`
+ SELECT id, indiko_sub, username, name, email, photo, url, role, created_at, updated_at
+ FROM users
+ WHERE id = ?
+ `, id).Scan(&user.ID, &user.IndikoSub, &user.Username, &user.Name, &user.Email, &user.Photo, &user.URL, &user.Role, &user.CreatedAt, &user.UpdatedAt)
+
+ if err == sql.ErrNoRows {
+ return nil, nil
+ }
+
+ if err != nil {
+ return nil, fmt.Errorf("query user: %w", err)
+ }
+
+ return user, nil
+}
+
+func (db *DB) UpdateUser(user *User) error {
+ user.UpdatedAt = time.Now().Unix()
+
+ _, err := db.Exec(`
+ UPDATE users
+ SET username = ?, name = ?, email = ?, photo = ?, url = ?, updated_at = ?
+ WHERE id = ?
+ `, user.Username, user.Name, user.Email, user.Photo, user.URL, user.UpdatedAt, user.ID)
+
+ if err != nil {
+ return fmt.Errorf("update user: %w", err)
+ }
+
+ return nil
+}
+
+func (db *DB) CreateSession(userID, accessToken, refreshToken string, expiresAt int64) (*Session, error) {
+ session := &Session{
+ ID: uuid.New().String(),
+ UserID: userID,
+ AccessToken: accessToken,
+ RefreshToken: refreshToken,
+ ExpiresAt: expiresAt,
+ CreatedAt: time.Now().Unix(),
+ }
+
+ _, err := db.Exec(`
+ INSERT INTO sessions (id, user_id, access_token, refresh_token, expires_at, created_at)
+ VALUES (?, ?, ?, ?, ?, ?)
+ `, session.ID, session.UserID, session.AccessToken, session.RefreshToken, session.ExpiresAt, session.CreatedAt)
+
+ if err != nil {
+ return nil, fmt.Errorf("insert session: %w", err)
+ }
+
+ return session, nil
+}
+
+func (db *DB) GetSessionByID(id string) (*Session, error) {
+ session := &Session{}
+ err := db.QueryRow(`
+ SELECT id, user_id, access_token, refresh_token, expires_at, created_at
+ FROM sessions
+ WHERE id = ?
+ `, id).Scan(&session.ID, &session.UserID, &session.AccessToken, &session.RefreshToken, &session.ExpiresAt, &session.CreatedAt)
+
+ if err == sql.ErrNoRows {
+ return nil, nil
+ }
+
+ if err != nil {
+ return nil, fmt.Errorf("query session: %w", err)
+ }
+
+ return session, nil
+}
+
+func (db *DB) DeleteSession(id string) error {
+ _, err := db.Exec("DELETE FROM sessions WHERE id = ?", id)
+ if err != nil {
+ return fmt.Errorf("delete session: %w", err)
+ }
+ return nil
+}
+
+func (db *DB) DeleteExpiredSessions() error {
+ now := time.Now().Unix()
+ _, err := db.Exec("DELETE FROM sessions WHERE expires_at < ?", now)
+ if err != nil {
+ return fmt.Errorf("delete expired sessions: %w", err)
+ }
+ return nil
+}
diff --git a/tsconfig.json b/tsconfig.json
deleted file mode 100644
index 1e1c5c9..0000000
--- a/tsconfig.json
+++ /dev/null
@@ -1,33 +0,0 @@
-{
- "compilerOptions": {
- // Environment setup & latest features
- "lib": ["ESNext", "DOM", "DOM.Iterable"],
- "target": "ESNext",
- "module": "Preserve",
- "moduleDetection": "force",
- "jsx": "preserve",
- "allowJs": true,
-
- // Bundler mode
- "moduleResolution": "bundler",
- "allowImportingTsExtensions": true,
- "verbatimModuleSyntax": true,
- "noEmit": true,
-
- // Decorators
- "experimentalDecorators": true,
- "useDefineForClassFields": false,
-
- // Best practices
- "strict": true,
- "skipLibCheck": true,
- "noFallthroughCasesInSwitch": true,
- "noUncheckedIndexedAccess": true,
- "noImplicitOverride": true,
-
- // Some stricter flags (disabled by default)
- "noUnusedLocals": false,
- "noUnusedParameters": false,
- "noPropertyAccessFromIndexSignature": false
- }
-}
diff --git a/web/server.go b/web/server.go
new file mode 100644
index 0000000..589109c
--- /dev/null
+++ b/web/server.go
@@ -0,0 +1,362 @@
+package web
+
+import (
+ "context"
+ "encoding/json"
+ "fmt"
+ "html/template"
+ "net/http"
+
+ "github.com/charmbracelet/log"
+ "github.com/kierank/pipes/auth"
+ "github.com/kierank/pipes/config"
+ "github.com/kierank/pipes/engine"
+ "github.com/kierank/pipes/store"
+)
+
+type Server struct {
+ cfg *config.Config
+ db *store.DB
+ server *http.Server
+ sessionManager *auth.SessionManager
+ oauthClient *auth.OAuthClient
+ templates *template.Template
+ logger *log.Logger
+}
+
+func NewServer(cfg *config.Config, db *store.DB, logger *log.Logger) *Server {
+ return &Server{
+ cfg: cfg,
+ db: db,
+ sessionManager: auth.NewSessionManager(cfg, db),
+ oauthClient: auth.NewOAuthClient(cfg, db),
+ logger: logger,
+ }
+}
+
+func (s *Server) Start() error {
+ // Load templates
+ tmpl, err := template.ParseGlob("web/templates/*.html")
+ if err != nil {
+ return fmt.Errorf("failed to load templates: %w", err)
+ }
+ s.templates = tmpl
+
+ mux := http.NewServeMux()
+
+ // Static files
+ mux.Handle("/public/", http.StripPrefix("/public/", http.FileServer(http.Dir("public"))))
+
+ // Public routes
+ mux.HandleFunc("/", s.handleIndex)
+ mux.HandleFunc("/health", s.handleHealth)
+
+ // Auth routes
+ mux.HandleFunc("/auth/login", s.handleLogin)
+ mux.HandleFunc("/auth/callback", s.handleCallback)
+ mux.HandleFunc("/auth/logout", s.handleLogout)
+
+ // Protected routes
+ mux.HandleFunc("/dashboard", s.sessionManager.RequireAuth(s.handleDashboard))
+ mux.HandleFunc("/pipes/", s.sessionManager.RequireAuth(s.handlePipeEditor))
+
+ // API routes
+ mux.HandleFunc("/api/me", s.sessionManager.RequireAuth(s.handleAPIMe))
+ mux.HandleFunc("/api/pipes", s.sessionManager.RequireAuth(s.handleAPIPipes))
+ mux.HandleFunc("/api/pipes/", s.sessionManager.RequireAuth(s.handleAPIPipe))
+ mux.HandleFunc("/api/node-types", s.handleAPINodeTypes)
+
+ s.server = &http.Server{
+ Addr: fmt.Sprintf("%s:%d", s.cfg.Host, s.cfg.Port),
+ Handler: mux,
+ }
+
+ return s.server.ListenAndServe()
+}
+
+func (s *Server) Shutdown(ctx context.Context) error {
+ if s.server != nil {
+ return s.server.Shutdown(ctx)
+ }
+ return nil
+}
+
+// Handlers
+
+func (s *Server) handleIndex(w http.ResponseWriter, r *http.Request) {
+ // Check if user is authenticated
+ user, _ := s.sessionManager.GetCurrentUser(r)
+ if user != nil {
+ http.Redirect(w, r, "/dashboard", http.StatusSeeOther)
+ return
+ }
+
+ w.Header().Set("Content-Type", "text/html")
+ s.templates.ExecuteTemplate(w, "index.html", nil)
+}
+
+func (s *Server) handleHealth(w http.ResponseWriter, r *http.Request) {
+ w.Write([]byte("OK"))
+}
+
+func (s *Server) handleLogin(w http.ResponseWriter, r *http.Request) {
+ authURL, err := s.oauthClient.GetAuthorizationURL()
+ if err != nil {
+ s.logger.Error("failed to generate auth URL", "error", err)
+ s.renderError(w, "Configuration Error", "Failed to start authentication process. Please contact the administrator.", err.Error())
+ return
+ }
+
+ http.Redirect(w, r, authURL, http.StatusSeeOther)
+}
+
+func (s *Server) handleCallback(w http.ResponseWriter, r *http.Request) {
+ code := r.URL.Query().Get("code")
+ state := r.URL.Query().Get("state")
+
+ if code == "" || state == "" {
+ s.renderError(w, "Invalid Request", "Missing authorization code or state parameter.", "")
+ return
+ }
+
+ user, session, err := s.oauthClient.HandleCallback(state, code)
+ if err != nil {
+ s.logger.Error("oauth callback error", "error", err)
+ s.renderError(w, "Authentication Failed", "We couldn't sign you in with Indiko. Please try again.", err.Error())
+ return
+ }
+
+ if err := s.sessionManager.SetSession(w, r, session.ID); err != nil {
+ s.logger.Error("failed to set session", "error", err)
+ s.renderError(w, "Session Error", "Authentication succeeded, but we couldn't create your session.", err.Error())
+ return
+ }
+
+ s.logger.Info("user authenticated", "name", user.Name, "email", user.Email)
+ http.Redirect(w, r, "/dashboard", http.StatusSeeOther)
+}
+
+func (s *Server) handleLogout(w http.ResponseWriter, r *http.Request) {
+ sessionID, _ := s.sessionManager.GetSessionID(r)
+ if sessionID != "" {
+ s.db.DeleteSession(sessionID)
+ }
+
+ s.sessionManager.ClearSession(w, r)
+ http.Redirect(w, r, "/", http.StatusSeeOther)
+}
+
+func (s *Server) handleDashboard(w http.ResponseWriter, r *http.Request) {
+ user := auth.GetUserFromContext(r.Context())
+ if user == nil {
+ http.Redirect(w, r, "/auth/login", http.StatusSeeOther)
+ return
+ }
+
+ pipes, err := s.db.GetUserPipes(user.ID)
+ if err != nil {
+ s.logger.Error("failed to get pipes", "user_id", user.ID, "error", err)
+ http.Error(w, "Failed to load pipes", http.StatusInternalServerError)
+ return
+ }
+
+ data := map[string]interface{}{
+ "User": user,
+ "Pipes": pipes,
+ }
+
+ w.Header().Set("Content-Type", "text/html")
+ s.templates.ExecuteTemplate(w, "dashboard.html", data)
+}
+
+func (s *Server) handlePipeEditor(w http.ResponseWriter, r *http.Request) {
+ // TODO: Implement pipe editor
+ w.Write([]byte("Pipe editor - coming soon!"))
+}
+
+func (s *Server) handleAPIMe(w http.ResponseWriter, r *http.Request) {
+ user := auth.GetUserFromContext(r.Context())
+ if user == nil {
+ http.Error(w, "Unauthorized", http.StatusUnauthorized)
+ return
+ }
+
+ w.Header().Set("Content-Type", "application/json")
+ json.NewEncoder(w).Encode(user)
+}
+
+func (s *Server) handleAPIPipes(w http.ResponseWriter, r *http.Request) {
+ user := auth.GetUserFromContext(r.Context())
+ if user == nil {
+ http.Error(w, "Unauthorized", http.StatusUnauthorized)
+ return
+ }
+
+ switch r.Method {
+ case "GET":
+ pipes, err := s.db.GetUserPipes(user.ID)
+ if err != nil {
+ http.Error(w, "Failed to load pipes", http.StatusInternalServerError)
+ return
+ }
+
+ w.Header().Set("Content-Type", "application/json")
+ json.NewEncoder(w).Encode(pipes)
+
+ case "POST":
+ var req struct {
+ Name string `json:"name"`
+ Description string `json:"description"`
+ Config string `json:"config"`
+ }
+
+ if err := json.NewDecoder(r.Body).Decode(&req); err != nil {
+ http.Error(w, "Invalid request", http.StatusBadRequest)
+ return
+ }
+
+ if req.Config == "" {
+ req.Config = `{"version":"1","nodes":[],"connections":[],"settings":{"enabled":false}}`
+ }
+
+ pipe, err := s.db.CreatePipe(user.ID, req.Name, req.Description, req.Config, false)
+ if err != nil {
+ http.Error(w, "Failed to create pipe", http.StatusInternalServerError)
+ return
+ }
+
+ w.Header().Set("Content-Type", "application/json")
+ w.WriteHeader(http.StatusCreated)
+ json.NewEncoder(w).Encode(pipe)
+
+ default:
+ http.Error(w, "Method not allowed", http.StatusMethodNotAllowed)
+ }
+}
+
+func (s *Server) handleAPIPipe(w http.ResponseWriter, r *http.Request) {
+ user := auth.GetUserFromContext(r.Context())
+ if user == nil {
+ http.Error(w, "Unauthorized", http.StatusUnauthorized)
+ return
+ }
+
+ // Extract pipe ID from path
+ pipeID := r.URL.Path[len("/api/pipes/"):]
+
+ switch r.Method {
+ case "GET":
+ pipe, err := s.db.GetPipe(pipeID)
+ if err != nil || pipe == nil {
+ http.Error(w, "Pipe not found", http.StatusNotFound)
+ return
+ }
+
+ if pipe.UserID != user.ID && !pipe.IsPublic {
+ http.Error(w, "Forbidden", http.StatusForbidden)
+ return
+ }
+
+ w.Header().Set("Content-Type", "application/json")
+ json.NewEncoder(w).Encode(pipe)
+
+ case "PUT":
+ pipe, err := s.db.GetPipe(pipeID)
+ if err != nil || pipe == nil {
+ http.Error(w, "Pipe not found", http.StatusNotFound)
+ return
+ }
+
+ if pipe.UserID != user.ID {
+ http.Error(w, "Forbidden", http.StatusForbidden)
+ return
+ }
+
+ var req struct {
+ Name string `json:"name"`
+ Description string `json:"description"`
+ Config map[string]interface{} `json:"config"`
+ }
+
+ if err := json.NewDecoder(r.Body).Decode(&req); err != nil {
+ http.Error(w, "Invalid request", http.StatusBadRequest)
+ return
+ }
+
+ if req.Name != "" {
+ pipe.Name = req.Name
+ }
+ if req.Description != "" {
+ pipe.Description = req.Description
+ }
+ if req.Config != nil {
+ configJSON, _ := json.Marshal(req.Config)
+ pipe.Config = string(configJSON)
+ }
+
+ if err := s.db.UpdatePipe(pipe); err != nil {
+ http.Error(w, "Failed to update pipe", http.StatusInternalServerError)
+ return
+ }
+
+ w.Header().Set("Content-Type", "application/json")
+ json.NewEncoder(w).Encode(map[string]bool{"success": true})
+
+ case "DELETE":
+ pipe, err := s.db.GetPipe(pipeID)
+ if err != nil || pipe == nil {
+ http.Error(w, "Pipe not found", http.StatusNotFound)
+ return
+ }
+
+ if pipe.UserID != user.ID {
+ http.Error(w, "Forbidden", http.StatusForbidden)
+ return
+ }
+
+ if err := s.db.DeletePipe(pipeID); err != nil {
+ http.Error(w, "Failed to delete pipe", http.StatusInternalServerError)
+ return
+ }
+
+ w.Header().Set("Content-Type", "application/json")
+ json.NewEncoder(w).Encode(map[string]bool{"success": true})
+
+ default:
+ http.Error(w, "Method not allowed", http.StatusMethodNotAllowed)
+ }
+}
+
+func (s *Server) handleAPINodeTypes(w http.ResponseWriter, r *http.Request) {
+ registry := engine.NewRegistry()
+ nodes := registry.GetAll()
+
+ var nodeTypes []map[string]interface{}
+ for _, node := range nodes {
+ nodeTypes = append(nodeTypes, map[string]interface{}{
+ "type": node.Type(),
+ "label": node.Label(),
+ "description": node.Description(),
+ "category": node.Category(),
+ "schema": node.GetConfigSchema(),
+ })
+ }
+
+ w.Header().Set("Content-Type", "application/json")
+ json.NewEncoder(w).Encode(nodeTypes)
+}
+
+// Helper functions
+
+func (s *Server) renderError(w http.ResponseWriter, title, message, details string) {
+ w.Header().Set("Content-Type", "text/html")
+ w.WriteHeader(http.StatusBadRequest)
+
+ data := map[string]interface{}{
+ "Title": title,
+ "Message": message,
+ "Details": details,
+ }
+
+ s.templates.ExecuteTemplate(w, "error.html", data)
+}
diff --git a/web/templates/dashboard.html b/web/templates/dashboard.html
new file mode 100644
index 0000000..9abc242
--- /dev/null
+++ b/web/templates/dashboard.html
@@ -0,0 +1,194 @@
+
+
+
+
+
+ Dashboard - Pipes
+
+
+
+
+
+
+
+
+
+
+
+
Your Pipes
+ {{if .Pipes}}
+
+ {{range .Pipes}}
+
+
{{.Name}}
+ {{if .Description}}
+
{{.Description}}
+ {{end}}
+
Edit
+
+ {{end}}
+
+ {{else}}
+
+
You haven't created any pipes yet!
+
+
+ {{end}}
+
+
+
+
+
+
diff --git a/web/templates/error.html b/web/templates/error.html
new file mode 100644
index 0000000..9d9922f
--- /dev/null
+++ b/web/templates/error.html
@@ -0,0 +1,112 @@
+
+
+
+
+
+ Error - Pipes
+
+
+
+
+
+
+
+
+
Error
+
{{.Title}}
+
{{.Message}}
+ {{if .Details}}
+
{{.Details}}
+ {{end}}
+
+
+
+
diff --git a/web/templates/index.html b/web/templates/index.html
new file mode 100644
index 0000000..e957dfc
--- /dev/null
+++ b/web/templates/index.html
@@ -0,0 +1,79 @@
+
+
+
+
+
+ Pipes - Visual Data Pipeline Builder
+
+
+
+
+
+
+
+
+
+