From 6969a0acad6532cbce2684fdfadab5fc4ef465b2 Mon Sep 17 00:00:00 2001 From: Kieran Klukas Date: Sat, 10 Jan 2026 22:17:43 -0500 Subject: [PATCH] feat: refactor to go and fix indiko --- .env.example | 16 +- .gitignore | 23 ++- CLAUDE.md | 72 +++++++ CRUSH.md | 187 ++++++++++++++++++ README.md | 212 +++++++++++++++++--- auth/middleware.go | 31 +++ auth/oauth.go | 242 +++++++++++++++++++++++ auth/session.go | 77 ++++++++ config.yaml.example | 23 +++ config/app.go | 161 ++++++++++++++++ engine/executor.go | 218 +++++++++++++++++++++ engine/registry.go | 59 ++++++ engine/scheduler.go | 93 +++++++++ go.mod | 40 ++++ go.sum | 84 ++++++++ main.go | 243 +++++++++++++++++++++++ nodes/node.go | 61 ++++++ nodes/sources/rss.go | 93 +++++++++ nodes/transforms/filter.go | 123 ++++++++++++ nodes/transforms/limit.go | 54 ++++++ nodes/transforms/sort.go | 91 +++++++++ package.json | 23 --- src/index.ts | 54 ------ src/pages/index.html | 15 -- src/types/env.d.ts | 7 - store/db.go | 148 ++++++++++++++ store/executions.go | 221 +++++++++++++++++++++ store/pipes.go | 212 ++++++++++++++++++++ store/users.go | 170 ++++++++++++++++ tsconfig.json | 33 ---- web/server.go | 362 +++++++++++++++++++++++++++++++++++ web/templates/dashboard.html | 194 +++++++++++++++++++ web/templates/error.html | 112 +++++++++++ web/templates/index.html | 79 ++++++++ 34 files changed, 3663 insertions(+), 170 deletions(-) create mode 100644 CLAUDE.md create mode 100644 CRUSH.md create mode 100644 auth/middleware.go create mode 100644 auth/oauth.go create mode 100644 auth/session.go create mode 100644 config.yaml.example create mode 100644 config/app.go create mode 100644 engine/executor.go create mode 100644 engine/registry.go create mode 100644 engine/scheduler.go create mode 100644 go.mod create mode 100644 go.sum create mode 100644 main.go create mode 100644 nodes/node.go create mode 100644 nodes/sources/rss.go create mode 100644 nodes/transforms/filter.go create mode 100644 nodes/transforms/limit.go create mode 100644 nodes/transforms/sort.go delete mode 100644 package.json delete mode 100644 src/index.ts delete mode 100644 src/pages/index.html delete mode 100644 src/types/env.d.ts create mode 100644 store/db.go create mode 100644 store/executions.go create mode 100644 store/pipes.go create mode 100644 store/users.go delete mode 100644 tsconfig.json create mode 100644 web/server.go create mode 100644 web/templates/dashboard.html create mode 100644 web/templates/error.html create mode 100644 web/templates/index.html 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 + + + + + + + +
+
+

Pipes Dashboard

+ +
+ +
+

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 + + + + + + + +
+

Pipes

+

A visual data pipeline builder inspired by Yahoo Pipes

+ Sign in with Indiko +
+ + -- 2.51.2