Mercurial
changeset 230:7795e3149540
[merge] Record closed hg-web feature branch
| author | MrJuneJune <me@mrjunejune.com> |
|---|---|
| date | Sun, 02 Aug 2026 16:48:18 -0700 |
| parents | 62f9cb40d10d (diff) 2de49ee4fdd6 (current diff) |
| children | d27054a14579 |
| files | |
| diffstat | 18 files changed, 1853 insertions(+), 34 deletions(-) [+] |
line wrap: on
line diff
--- a/.bazelrc Sun Aug 02 16:48:03 2026 -0700 +++ b/.bazelrc Sun Aug 02 16:48:18 2026 -0700 @@ -2,7 +2,8 @@ common --enable_platform_specific_config # Suppress duplicate library warnings from openssl BCR module on macOS -build:macos --linkopt=-Wl,-no_warn_duplicate_libraries +# Use host_linkopt to avoid passing this to cross-compilation targets (e.g., WASM) +build:macos --host_linkopt=-Wl,-no_warn_duplicate_libraries # esmdk build --incompatible_enable_cc_toolchain_resolution
--- a/MODULE.bazel Sun Aug 02 16:48:03 2026 -0700 +++ b/MODULE.bazel Sun Aug 02 16:48:18 2026 -0700 @@ -63,8 +63,10 @@ sha256 = "cf0ed0a920799d576ffde4e0cae66d732bf23c2530407f26f59c7831dffe1f0e", ) +# Bring in Python support +bazel_dep(name = "rules_python", version = "1.7.0") + # Bring in pip support -# bazel_dep(name = "rules_python", version = "1.7.0") # use_extension("@rules_python//python/extensions:pip.bzl", "pip") # pip = use_extension("@rules_python//python/extensions:pip.bzl", "pip")
--- a/hg-web/BUILD Sun Aug 02 16:48:03 2026 -0700 +++ b/hg-web/BUILD Sun Aug 02 16:48:18 2026 -0700 @@ -87,3 +87,12 @@ name = "hg_web_server_bundle", binary = ":hg_web_server", ) + +sh_binary( + name = "deploy", + srcs = ["deploy.sh"], + data = [ + ":hg_web_server_bundle", + "@bazel_tools//tools/bash/runfiles", + ], +)
--- a/hg-web/README.md Sun Aug 02 16:48:03 2026 -0700 +++ b/hg-web/README.md Sun Aug 02 16:48:18 2026 -0700 @@ -217,6 +217,25 @@ release root, active path, health URL, user, and group can be overridden with environment variables. +## Publishing + +Push committed Mercurial changes to the configured server: + +```bash +hg push +``` + +Then update and deploy on the server: + +```bash +ssh -t [email protected] \ + 'cd ~/zenbu && hg update default && bazel run //hg-web:deploy' +``` + +The SSH key may prompt for its passphrase. `//hg-web:deploy` builds the release +bundle through Bazel before running the atomic promotion, health check, and +rollback workflow. + ## Forge capability tree ```text
--- a/hg-web/deploy.sh Sun Aug 02 16:48:03 2026 -0700 +++ b/hg-web/deploy.sh Sun Aug 02 16:48:18 2026 -0700 @@ -8,15 +8,45 @@ SERVICE_USER="${SERVICE_USER:-hg_web_server}" SERVICE_GROUP="${SERVICE_GROUP:-zenbu_team}" -workspace="$(hg root)" +workspace="${BUILD_WORKSPACE_DIRECTORY:-$(hg root)}" cd "$workspace" +if [[ -n "${BUILD_WORKSPACE_DIRECTORY:-}" ]]; then + if [[ -n "${RUNFILES_DIR:-}" ]]; then + source "$RUNFILES_DIR/bazel_tools/tools/bash/runfiles/runfiles.bash" + elif [[ -n "${RUNFILES_MANIFEST_FILE:-}" ]]; then + runfiles_library="$( + grep -sm1 '^bazel_tools/tools/bash/runfiles/runfiles.bash ' \ + "$RUNFILES_MANIFEST_FILE" | cut -d' ' -f2- + )" + source "$runfiles_library" + elif [[ -d "$0.runfiles" ]]; then + RUNFILES_DIR="$0.runfiles" + export RUNFILES_DIR + source "$RUNFILES_DIR/bazel_tools/tools/bash/runfiles/runfiles.bash" + elif [[ -f "$0.runfiles_manifest" ]]; then + RUNFILES_MANIFEST_FILE="$0.runfiles_manifest" + export RUNFILES_MANIFEST_FILE + runfiles_library="$( + grep -sm1 '^bazel_tools/tools/bash/runfiles/runfiles.bash ' \ + "$RUNFILES_MANIFEST_FILE" | cut -d' ' -f2- + )" + source "$runfiles_library" + else + echo "Bazel runfiles are unavailable." >&2 + exit 1 + fi + bundle_dir="$(rlocation _main/hg-web/hg_web_server_bundle)" +else + bazel build -c opt //hg-web:hg_web_server_bundle + bundle_dir="bazel-bin/hg-web/hg_web_server_bundle" +fi + revision="$(hg log -r . -T '{node|short}')" release_name="${revision}-$(date -u +%Y%m%dT%H%M%SZ)" release_dir="${RELEASE_ROOT}/${release_name}" staging_dir="${RELEASE_ROOT}/.${release_name}.tmp" next_link="${ACTIVE_PATH}.next" -bundle_dir="bazel-bin/hg-web/hg_web_server_bundle" previous_release="" promoted=0 @@ -53,8 +83,6 @@ } trap rollback ERR -bazel build -c opt //hg-web:hg_web_server_bundle - sudo install -d -o root -g "$SERVICE_GROUP" -m 0755 "$RELEASE_ROOT" sudo rm -rf "$staging_dir" sudo install -d -o "$SERVICE_USER" -g "$SERVICE_GROUP" -m 0755 "$staging_dir"
--- a/load_test/README.md Sun Aug 02 16:48:03 2026 -0700 +++ /dev/null Thu Jan 01 00:00:00 1970 +0000 @@ -1,7 +0,0 @@ -# load_test - -Load testing and performance measurement scripts. - -## Files - -- `main.py` - Python load testing script
--- a/load_test/main.py Sun Aug 02 16:48:03 2026 -0700 +++ /dev/null Thu Jan 01 00:00:00 1970 +0000 @@ -1,8 +0,0 @@ -from locust import HttpUser, task, between - -class WebsiteUser(HttpUser): - wait_time = between(1, 5) - - @task(3) - def index_page(self): - self.client.get("/")
--- a/mrjunejune/main.c Sun Aug 02 16:48:03 2026 -0700 +++ b/mrjunejune/main.c Sun Aug 02 16:48:18 2026 -0700 @@ -25,6 +25,11 @@ S3_Config s3_config; } Media_Processing_Context; +typedef struct { + char *input_path; + char *output_path; +} File_Converter_Config; + // Server configuration (loaded from .config) static char g_upload_auth_token[256] = {0}; static char g_s3_region[64] = "us-west-2"; @@ -290,6 +295,27 @@ return resp; } +// Background thread function for media processing +void *Simple_WebpConverter_Background(void *arg) +{ + File_Converter_Config *configuration = (File_Converter_Config *)arg; + + char cmd[1024]; + snprintf(cmd, sizeof(cmd), "ffmpeg -y -i %s -quality 80 %s 2>/tmp/error_log", + configuration->input_path, configuration->output_path); + Seobeo_Log(SEOBEO_INFO, "[MEDIA] Running FFmpeg: %s\n", cmd); + int ffmpeg_result = system(cmd); + + Seobeo_Log(SEOBEO_INFO, "[MEDIA] FFmpeg result: %d\n", ffmpeg_result); + if (ffmpeg_result != 0) + { + Seobeo_Log(SEOBEO_ERROR, "[MEDIA] ERROR: FFmpeg conversion failed\n"); + return NULL; + } + Seobeo_Log(SEOBEO_INFO, "[MEDIA] Successfully converted to webp: %s\n"); + return NULL; +} + Seobeo_Request_Entry *ConvertImageToWebP(Seobeo_Request_Entry *req, Dowa_Arena *arena) { Seobeo_Request_Entry *resp = NULL; @@ -338,7 +364,7 @@ const char *content_length_str = ((Seobeo_Request_Entry*)cl_kv)->value; size_t file_size = atoi(content_length_str); - printf("DEBUG: Converting image, file_size=%zu bytes\n", file_size); + Seobeo_Log(SEOBEO_DEBUG, "Converting image, file_size=%zu bytes\n", file_size); int open_flags = O_RDWR | O_CREAT | O_EXCL; @@ -365,15 +391,15 @@ Dowa_String_UUID(seed, uuid4); char *output_path = (char *)Dowa_Arena_Allocate(arena, TMP_FILE_LENGTH);; snprintf(output_path, TMP_FILE_LENGTH, "/tmp/%s.webp", uuid4); - printf("[DEBUG] output_path %s\n", output_path); - printf("[DEBUG] open_flags: 0x%x\n", open_flags); - printf("[DEBUG] input_path: %s\n", input_path); + Seobeo_Log(SEOBEO_DEBUG, "output_path %s\n", output_path); + Seobeo_Log(SEOBEO_DEBUG, "open_flags: 0x%x\n", open_flags); + Seobeo_Log(SEOBEO_DEBUG, "input_path: %s\n", input_path); int output_fd = open(output_path, open_flags, 0600); - printf("[DEBUG] output_fd: %d\n", output_fd); + Seobeo_Log(SEOBEO_DEBUG, "output_fd: %d\n", output_fd); if (output_fd == -1) { unlink(input_path); - printf("[DEBUG] errno: %d (%s)\n", errno, strerror(errno)); + Seobeo_Log(SEOBEO_DEBUG, "errno: %d (%s)\n", errno, strerror(errno)); char *error_msg = "Failed to create output file"; Dowa_HashMap_Push_Arena(resp, "status", "500", arena); Dowa_HashMap_Push_Arena(resp, "content-type", "text/plain", arena); @@ -382,11 +408,14 @@ } close(output_fd); - char cmd[1024]; - snprintf(cmd, sizeof(cmd), "ffmpeg -y -i %s -quality 80 %s 2>/tmp/error_log", - input_path, output_path); - int result = system(cmd); - if (result != 0) + File_Converter_Config *configuration = Dowa_Arena_Allocate(arena, sizeof(File_Converter_Config)); + configuration->input_path = input_path; + configuration->output_path = output_path; + + pthread_t thread_id; + int thread_result = pthread_create(&thread_id, NULL, Simple_WebpConverter_Background, (void *)configuration); + + if (thread_result != 0) { unlink(input_path); unlink(output_path); @@ -396,6 +425,12 @@ Dowa_HashMap_Push_Arena(resp, "body", error_msg, arena); return resp; } + else + { + // Detach thread so it cleans up automatically when done + pthread_detach(thread_id); + Seobeo_Log(SEOBEO_INFO, "[MEDIA] Successfully spawned and detached thread\n"); + } size_t converted_size = 0; FILE *out_file = fopen(output_path, "rb"); @@ -412,7 +447,6 @@ fclose(out_file); unlink(input_path); - char *filename = strrchr(output_path, '/') + 1; char *response_body = Dowa_Arena_Allocate(arena, 512); snprintf(response_body, 512, @@ -422,7 +456,7 @@ Dowa_HashMap_Push_Arena(resp, "content-type", "application/json", arena); Dowa_HashMap_Push_Arena(resp, "body", response_body, arena); - printf("DEBUG: Image converted, available at /api/download/%s\n", filename); + Seobeo_Log(SEOBEO_DEBUG, "Image converted, available at /api/download/%s\n", filename); return resp; }
--- /dev/null Thu Jan 01 00:00:00 1970 +0000 +++ b/schwab_trader/BUILD Sun Aug 02 16:48:18 2026 -0700 @@ -0,0 +1,66 @@ +load("@rules_python//python:py_binary.bzl", "py_binary") +load("@rules_python//python:py_library.bzl", "py_library") +load("@rules_python//python:py_test.bzl", "py_test") + +py_library( + name = "schwab_client", + srcs = [ + "__init__.py", + "schwab_client.py", + ], + visibility = ["//visibility:public"], +) + +py_library( + name = "schwab_dashboard_lib", + srcs = [ + "__init__.py", + "dashboard.py", + ], + deps = [":schwab_client"], +) + +py_library( + name = "schwab_dashboard_server_lib", + srcs = [ + "__init__.py", + "dashboard_server.py", + ], + deps = [ + ":schwab_client", + ":schwab_dashboard_lib", + ], +) + +py_binary( + name = "schwab_cli", + srcs = ["schwab_cli.py"], + main = "schwab_cli.py", + deps = [":schwab_client"], +) + +py_binary( + name = "schwab_dashboard", + srcs = ["dashboard_server.py"], + main = "dashboard_server.py", + deps = [ + ":schwab_client", + ":schwab_dashboard_lib", + ], +) + +py_test( + name = "schwab_client_test", + srcs = ["schwab_client_test.py"], + deps = [":schwab_client"], +) + +py_test( + name = "dashboard_test", + srcs = ["dashboard_test.py"], + deps = [ + ":schwab_client", + ":schwab_dashboard_lib", + ":schwab_dashboard_server_lib", + ], +)
--- /dev/null Thu Jan 01 00:00:00 1970 +0000 +++ b/schwab_trader/README.md Sun Aug 02 16:48:18 2026 -0700 @@ -0,0 +1,160 @@ +# Schwab Trader + +Small Bazel-built helper for the Schwab Trader API. This project is only plumbing: it helps authenticate, inspect accounts, build explicit user-specified stock orders, and submit them only when a live-trade confirmation flag is present. + +It does not recommend trades, choose symbols, allocate portfolio risk, or automate a strategy. + +## Can Schwab accounts be traded by API? + +Yes. Schwab provides the Trader API through the Schwab Developer Portal. The flow is OAuth 2.0: + +1. Create an app in the Schwab Developer Portal. +2. Get an app key and app secret. +3. Register a redirect URI. +4. Open the OAuth authorization URL and log in through Schwab. +5. Exchange the returned `code` for an access token and refresh token. +6. Use the access token against `https://api.schwabapi.com/trader/v1`. + +Do not use username/password scraping or browser automation. The supported path is OAuth tokens. + +Useful endpoints this project targets: + +- `GET https://api.schwabapi.com/trader/v1/accounts/accountNumbers` +- `GET https://api.schwabapi.com/trader/v1/accounts` +- `GET https://api.schwabapi.com/trader/v1/accounts/{accountHash}` +- `POST https://api.schwabapi.com/trader/v1/accounts/{accountHash}/orders` + +Schwab uses account hashes for trading API calls. Fetch them with `account-numbers` before placing any order. + +## Build and test + +From the repo root: + +```bash +bazel build //schwab_trader:schwab_cli +bazel test //schwab_trader:schwab_client_test +bazel build //schwab_trader:schwab_dashboard +bazel test //schwab_trader:dashboard_test +``` + +## Configuration + +Set these environment variables: + +```bash +export SCHWAB_APP_KEY="your-schwab-app-key" +export SCHWAB_APP_SECRET="your-schwab-app-secret" +export SCHWAB_REDIRECT_URI="https://127.0.0.1" +export SCHWAB_TOKEN_FILE="$HOME/.config/zenbu/schwab_tokens.json" +``` + +`SCHWAB_TOKEN_FILE` is optional and defaults to `~/.config/zenbu/schwab_tokens.json`. Token files are written with `0600` permissions. + +## OAuth bootstrap + +Print the Schwab login URL: + +```bash +bazel run //schwab_trader:schwab_cli -- auth-url +``` + +Open it, log in through Schwab, authorize the app, then copy the full callback URL or just its `code` parameter: + +```bash +bazel run //schwab_trader:schwab_cli -- token --code 'https://127.0.0.1/?code=...' +``` + +Refresh later: + +```bash +bazel run //schwab_trader:schwab_cli -- refresh +``` + +## Account discovery + +```bash +bazel run //schwab_trader:schwab_cli -- account-numbers +bazel run //schwab_trader:schwab_cli -- accounts --positions +``` + +## Order dry run + +Build a stock order payload without sending it: + +```bash +bazel run //schwab_trader:schwab_cli -- build-equity-order \ + --action BUY \ + --symbol AAPL \ + --quantity 1 \ + --order-type MARKET +``` + +`place-equity-order` is also dry-run by default: + +```bash +bazel run //schwab_trader:schwab_cli -- place-equity-order \ + --account-hash "$SCHWAB_ACCOUNT_HASH" \ + --action SELL \ + --symbol AAPL \ + --quantity 1 \ + --order-type LIMIT \ + --price 250.00 +``` + +To actually submit an order, both safety flags are required: + +```bash +bazel run //schwab_trader:schwab_cli -- place-equity-order \ + --account-hash "$SCHWAB_ACCOUNT_HASH" \ + --action BUY \ + --symbol AAPL \ + --quantity 1 \ + --order-type MARKET \ + --live \ + --confirm-live-trade +``` + +Use live trading only after checking the generated JSON, account hash, symbol, quantity, order type, and Schwab API permissions. + +## Local sentiment dashboard + +Run the local dashboard: + +```bash +bazel run //schwab_trader:schwab_dashboard +``` + +Then open: + +```text +http://127.0.0.1:8765 +``` + +The dashboard is intentionally local-first and read/paper-trade oriented. It has no live-trade endpoint. It shows: + +- Schwab environment/token status without exposing token values +- risk settings such as profit target, stop loss, confidence threshold, and max paper-trade dollars +- manually added social evidence from Reddit/X/news/etc. +- deterministic sentiment signals and confidence +- paper trades with simple risk rejection +- audit events + +The dashboard stores state in: + +```text +~/.local/share/zenbu/schwab_trader/dashboard.db +``` + +Override it when testing: + +```bash +bazel run //schwab_trader:schwab_dashboard -- --db /tmp/schwab_dashboard.db +``` + +Social evidence can be added through the page or API: + +```bash +curl -X POST http://127.0.0.1:8765/api/evidence \ + -H 'Content-Type: application/json' \ + -d '{"source":"reddit","symbol":"AAPL","text":"$AAPL bullish strong growth","engagement":42}' +```
--- /dev/null Thu Jan 01 00:00:00 1970 +0000 +++ b/schwab_trader/__init__.py Sun Aug 02 16:48:18 2026 -0700 @@ -0,0 +1,2 @@ +"""Schwab Trader API helpers for Zenbu.""" +
--- /dev/null Thu Jan 01 00:00:00 1970 +0000 +++ b/schwab_trader/dashboard.py Sun Aug 02 16:48:18 2026 -0700 @@ -0,0 +1,598 @@ +from __future__ import annotations + +import json +import os +import re +import sqlite3 +import time +from dataclasses import dataclass +from pathlib import Path +from typing import Any + +from schwab_trader.schwab_client import DEFAULT_TOKEN_FILE, SchwabConfig, load_tokens + + +DEFAULT_DASHBOARD_DB = "~/.local/share/zenbu/schwab_trader/dashboard.db" + +DEFAULT_SETTINGS: dict[str, Any] = { + "profit_target_pct": 3.0, + "stop_loss_pct": 2.0, + "max_trade_dollars": 500.0, + "max_account_pct": 2.0, + "max_open_positions": 3, + "min_confidence": 0.60, + "min_evidence_count": 3, + "min_source_count": 2, + "cooldown_minutes": 60, + "market_hours_only": True, + "max_daily_loss_dollars": 250.0, + "day_trade_limit": 3, + "live_trading_enabled": False, + "require_manual_confirmation": True, + "llm_enabled": False, + "allowlist": [], + "blocklist": [], +} + +POSITIVE_WORDS = { + "beat", + "beats", + "bull", + "bullish", + "buy", + "calls", + "growth", + "hype", + "moon", + "mooning", + "profit", + "rally", + "strong", + "surge", + "up", + "winner", +} + +NEGATIVE_WORDS = { + "bear", + "bearish", + "crash", + "dump", + "fall", + "falling", + "fraud", + "lawsuit", + "loss", + "miss", + "puts", + "risk", + "sell", + "short", + "weak", +} + +SYMBOL_RE = re.compile(r"(?<![A-Z0-9])\$?([A-Z]{1,5})(?![A-Z0-9])") +COMMON_WORDS = { + "A", + "AI", + "AM", + "API", + "CEO", + "CFO", + "DD", + "ETF", + "GDP", + "IPO", + "IRS", + "LLM", + "PDT", + "SEC", + "USA", + "USD", +} + + +@dataclass(frozen=True) +class EvidenceInput: + source: str + text: str + symbol: str | None = None + url: str | None = None + engagement: float = 0.0 + raw: dict[str, Any] | None = None + + +class DashboardStore: + def __init__(self, db_path: Path | str | None = None) -> None: + if db_path is None: + db_path = os.environ.get("SCHWAB_DASHBOARD_DB", DEFAULT_DASHBOARD_DB) + self.db_path = Path(db_path).expanduser() + self.db_path.parent.mkdir(parents=True, exist_ok=True) + self._init_db() + + def get_status(self) -> dict[str, Any]: + token_file = Path(os.environ.get("SCHWAB_TOKEN_FILE", DEFAULT_TOKEN_FILE)).expanduser() + token_status: dict[str, Any] = { + "path": str(token_file), + "exists": token_file.exists(), + "access_token_present": False, + "refresh_token_present": False, + "saved_at": None, + "age_seconds": None, + } + if token_file.exists(): + try: + tokens = load_tokens(token_file) + saved_at = tokens.get("saved_at") + token_status.update( + { + "access_token_present": bool(tokens.get("access_token")), + "refresh_token_present": bool(tokens.get("refresh_token")), + "saved_at": saved_at, + "age_seconds": int(time.time()) - int(saved_at) if saved_at else None, + } + ) + except (OSError, ValueError, TypeError) as error: + token_status["error"] = str(error) + + env_status = { + "SCHWAB_APP_KEY": bool(os.environ.get("SCHWAB_APP_KEY")), + "SCHWAB_APP_SECRET": bool(os.environ.get("SCHWAB_APP_SECRET")), + "SCHWAB_REDIRECT_URI": bool(os.environ.get("SCHWAB_REDIRECT_URI")), + } + + return { + "service": "schwab-dashboard", + "database": str(self.db_path), + "env": env_status, + "tokens": token_status, + "live_trading_enabled": False, + "live_trading_note": "Dashboard has no live-trade endpoint; use CLI dry-run/manual confirmation flow.", + "counts": self.get_counts(), + } + + def get_counts(self) -> dict[str, int]: + with self._connect() as conn: + return { + "evidence": self._count(conn, "evidence"), + "signals": self._count(conn, "signals"), + "paper_trades": self._count(conn, "paper_trades"), + "audit_events": self._count(conn, "audit"), + } + + def get_settings(self) -> dict[str, Any]: + settings = dict(DEFAULT_SETTINGS) + with self._connect() as conn: + for row in conn.execute("SELECT key, value_json FROM settings"): + settings[row["key"]] = json.loads(row["value_json"]) + return settings + + def update_settings(self, updates: dict[str, Any]) -> dict[str, Any]: + allowed = set(DEFAULT_SETTINGS) + unknown = sorted(set(updates) - allowed) + if unknown: + raise ValueError("Unknown settings: " + ", ".join(unknown)) + + current = self.get_settings() + current.update(updates) + self._validate_settings(current) + + with self._connect() as conn: + for key, value in current.items(): + conn.execute( + """ + INSERT INTO settings(key, value_json) + VALUES (?, ?) + ON CONFLICT(key) DO UPDATE SET value_json = excluded.value_json + """, + (key, json.dumps(value, sort_keys=True)), + ) + conn.commit() + self.add_audit("settings.updated", "Dashboard settings updated", updates) + return current + + def add_evidence(self, item: EvidenceInput) -> dict[str, Any]: + source = item.source.strip().lower() + text = item.text.strip() + if not source: + raise ValueError("source is required") + if not text: + raise ValueError("text is required") + + symbol = normalize_symbol(item.symbol) if item.symbol else extract_symbol(text) + if not symbol: + raise ValueError("symbol is required or must be detectable as a ticker in text") + + sentiment_score = score_sentiment(text) + created_at = int(time.time()) + + with self._connect() as conn: + cursor = conn.execute( + """ + INSERT INTO evidence( + source, symbol, url, text, engagement, sentiment_score, raw_json, created_at + ) + VALUES (?, ?, ?, ?, ?, ?, ?, ?) + """, + ( + source, + symbol, + item.url, + text, + float(item.engagement), + sentiment_score, + json.dumps(item.raw or {}, sort_keys=True), + created_at, + ), + ) + evidence_id = int(cursor.lastrowid) + conn.commit() + + signal = self.recompute_signal(symbol) + self.add_audit( + "evidence.added", + f"Added {source} evidence for {symbol}", + {"evidence_id": evidence_id, "symbol": symbol, "signal": signal}, + ) + return {"id": evidence_id, "symbol": symbol, "sentiment_score": sentiment_score, "signal": signal} + + def list_evidence(self, limit: int = 100) -> list[dict[str, Any]]: + with self._connect() as conn: + rows = conn.execute( + """ + SELECT id, source, symbol, url, text, engagement, sentiment_score, created_at + FROM evidence + ORDER BY id DESC + LIMIT ? + """, + (limit,), + ).fetchall() + return [dict(row) for row in rows] + + def recompute_signal(self, symbol: str) -> dict[str, Any]: + symbol = normalize_symbol(symbol) + settings = self.get_settings() + with self._connect() as conn: + rows = conn.execute( + """ + SELECT source, sentiment_score, engagement, created_at + FROM evidence + WHERE symbol = ? + ORDER BY id DESC + LIMIT 100 + """, + (symbol,), + ).fetchall() + + evidence_count = len(rows) + source_count = len({row["source"] for row in rows}) + weighted_total = 0.0 + weight_sum = 0.0 + for row in rows: + engagement_weight = min(5.0, 1.0 + max(0.0, float(row["engagement"])) / 100.0) + weighted_total += float(row["sentiment_score"]) * engagement_weight + weight_sum += engagement_weight + sentiment_score = weighted_total / weight_sum if weight_sum else 0.0 + confidence = compute_confidence(sentiment_score, evidence_count, source_count) + action = decide_signal_action(symbol, sentiment_score, confidence, evidence_count, source_count, settings) + summary = summarize_signal(symbol, sentiment_score, confidence, evidence_count, source_count, action) + updated_at = int(time.time()) + + conn.execute( + """ + INSERT INTO signals( + symbol, sentiment_score, confidence, evidence_count, source_count, + action, summary, updated_at + ) + VALUES (?, ?, ?, ?, ?, ?, ?, ?) + ON CONFLICT(symbol) DO UPDATE SET + sentiment_score = excluded.sentiment_score, + confidence = excluded.confidence, + evidence_count = excluded.evidence_count, + source_count = excluded.source_count, + action = excluded.action, + summary = excluded.summary, + updated_at = excluded.updated_at + """, + ( + symbol, + sentiment_score, + confidence, + evidence_count, + source_count, + action, + summary, + updated_at, + ), + ) + conn.commit() + + return { + "symbol": symbol, + "sentiment_score": round(sentiment_score, 4), + "confidence": round(confidence, 4), + "evidence_count": evidence_count, + "source_count": source_count, + "action": action, + "summary": summary, + "updated_at": updated_at, + } + + def list_signals(self) -> list[dict[str, Any]]: + with self._connect() as conn: + rows = conn.execute( + """ + SELECT symbol, sentiment_score, confidence, evidence_count, source_count, + action, summary, updated_at + FROM signals + ORDER BY confidence DESC, updated_at DESC + """ + ).fetchall() + return [dict(row) for row in rows] + + def add_paper_trade(self, payload: dict[str, Any]) -> dict[str, Any]: + symbol = normalize_symbol(str(payload.get("symbol", ""))) + action = str(payload.get("action", "")).upper() + quantity = float(payload.get("quantity", 0)) + price = float(payload.get("price", 0)) + reason = str(payload.get("reason", "manual paper trade")).strip() + + if action not in {"BUY", "SELL"}: + raise ValueError("action must be BUY or SELL") + if not symbol: + raise ValueError("symbol is required") + if quantity <= 0: + raise ValueError("quantity must be greater than zero") + if price <= 0: + raise ValueError("price must be greater than zero") + + settings = self.get_settings() + notional = quantity * price + status = "accepted" + if notional > float(settings["max_trade_dollars"]): + status = "rejected_max_trade_dollars" + elif symbol in {normalize_symbol(s) for s in settings["blocklist"]}: + status = "rejected_blocklist" + + created_at = int(time.time()) + with self._connect() as conn: + cursor = conn.execute( + """ + INSERT INTO paper_trades(symbol, action, quantity, price, notional, status, reason, created_at) + VALUES (?, ?, ?, ?, ?, ?, ?, ?) + """, + (symbol, action, quantity, price, notional, status, reason, created_at), + ) + trade_id = int(cursor.lastrowid) + conn.commit() + + result = { + "id": trade_id, + "symbol": symbol, + "action": action, + "quantity": quantity, + "price": price, + "notional": notional, + "status": status, + "reason": reason, + "created_at": created_at, + } + self.add_audit("paper_trade.created", f"Paper trade {status}: {action} {quantity} {symbol}", result) + return result + + def list_paper_trades(self, limit: int = 100) -> list[dict[str, Any]]: + with self._connect() as conn: + rows = conn.execute( + """ + SELECT id, symbol, action, quantity, price, notional, status, reason, created_at + FROM paper_trades + ORDER BY id DESC + LIMIT ? + """, + (limit,), + ).fetchall() + return [dict(row) for row in rows] + + def add_audit(self, event_type: str, message: str, payload: dict[str, Any] | None = None) -> None: + with self._connect() as conn: + conn.execute( + """ + INSERT INTO audit(event_type, message, payload_json, created_at) + VALUES (?, ?, ?, ?) + """, + (event_type, message, json.dumps(payload or {}, sort_keys=True), int(time.time())), + ) + conn.commit() + + def list_audit(self, limit: int = 100) -> list[dict[str, Any]]: + with self._connect() as conn: + rows = conn.execute( + """ + SELECT id, event_type, message, payload_json, created_at + FROM audit + ORDER BY id DESC + LIMIT ? + """, + (limit,), + ).fetchall() + events = [] + for row in rows: + event = dict(row) + event["payload"] = json.loads(event.pop("payload_json")) + events.append(event) + return events + + def _init_db(self) -> None: + with self._connect() as conn: + conn.executescript( + """ + CREATE TABLE IF NOT EXISTS settings ( + key TEXT PRIMARY KEY, + value_json TEXT NOT NULL + ); + + CREATE TABLE IF NOT EXISTS evidence ( + id INTEGER PRIMARY KEY AUTOINCREMENT, + source TEXT NOT NULL, + symbol TEXT NOT NULL, + url TEXT, + text TEXT NOT NULL, + engagement REAL NOT NULL DEFAULT 0, + sentiment_score REAL NOT NULL, + raw_json TEXT NOT NULL DEFAULT '{}', + created_at INTEGER NOT NULL + ); + + CREATE INDEX IF NOT EXISTS idx_evidence_symbol_created + ON evidence(symbol, created_at); + + CREATE TABLE IF NOT EXISTS signals ( + symbol TEXT PRIMARY KEY, + sentiment_score REAL NOT NULL, + confidence REAL NOT NULL, + evidence_count INTEGER NOT NULL, + source_count INTEGER NOT NULL, + action TEXT NOT NULL, + summary TEXT NOT NULL, + updated_at INTEGER NOT NULL + ); + + CREATE TABLE IF NOT EXISTS paper_trades ( + id INTEGER PRIMARY KEY AUTOINCREMENT, + symbol TEXT NOT NULL, + action TEXT NOT NULL, + quantity REAL NOT NULL, + price REAL NOT NULL, + notional REAL NOT NULL, + status TEXT NOT NULL, + reason TEXT NOT NULL, + created_at INTEGER NOT NULL + ); + + CREATE TABLE IF NOT EXISTS audit ( + id INTEGER PRIMARY KEY AUTOINCREMENT, + event_type TEXT NOT NULL, + message TEXT NOT NULL, + payload_json TEXT NOT NULL, + created_at INTEGER NOT NULL + ); + """ + ) + for key, value in DEFAULT_SETTINGS.items(): + conn.execute( + "INSERT OR IGNORE INTO settings(key, value_json) VALUES (?, ?)", + (key, json.dumps(value, sort_keys=True)), + ) + conn.commit() + + def _connect(self) -> sqlite3.Connection: + conn = sqlite3.connect(self.db_path) + conn.row_factory = sqlite3.Row + return conn + + @staticmethod + def _count(conn: sqlite3.Connection, table: str) -> int: + return int(conn.execute(f"SELECT COUNT(*) FROM {table}").fetchone()[0]) + + @staticmethod + def _validate_settings(settings: dict[str, Any]) -> None: + positive_numbers = [ + "profit_target_pct", + "stop_loss_pct", + "max_trade_dollars", + "max_account_pct", + "max_daily_loss_dollars", + ] + for key in positive_numbers: + if float(settings[key]) <= 0: + raise ValueError(f"{key} must be greater than zero") + if not 0 <= float(settings["min_confidence"]) <= 1: + raise ValueError("min_confidence must be between 0 and 1") + if int(settings["min_evidence_count"]) < 1: + raise ValueError("min_evidence_count must be at least 1") + if int(settings["min_source_count"]) < 1: + raise ValueError("min_source_count must be at least 1") + if bool(settings["live_trading_enabled"]): + raise ValueError("live_trading_enabled cannot be enabled from the dashboard") + + +def normalize_symbol(value: str) -> str: + symbol = value.strip().upper().lstrip("$") + if not re.fullmatch(r"[A-Z]{1,5}", symbol): + raise ValueError("symbol must be 1-5 letters") + return symbol + + +def extract_symbol(text: str) -> str | None: + for match in SYMBOL_RE.finditer(text.upper()): + symbol = match.group(1) + if symbol not in COMMON_WORDS: + return symbol + return None + + +def score_sentiment(text: str) -> float: + words = re.findall(r"[a-zA-Z']+", text.lower()) + positive = sum(1 for word in words if word in POSITIVE_WORDS) + negative = sum(1 for word in words if word in NEGATIVE_WORDS) + total = positive + negative + if total == 0: + return 0.0 + return max(-1.0, min(1.0, (positive - negative) / total)) + + +def compute_confidence(sentiment_score: float, evidence_count: int, source_count: int) -> float: + evidence_component = min(0.35, evidence_count * 0.07) + source_component = min(0.25, source_count * 0.10) + sentiment_component = min(0.20, abs(sentiment_score) * 0.20) + return min(0.95, 0.20 + evidence_component + source_component + sentiment_component) + + +def decide_signal_action( + symbol: str, + sentiment_score: float, + confidence: float, + evidence_count: int, + source_count: int, + settings: dict[str, Any], +) -> str: + blocklist = {normalize_symbol(item) for item in settings.get("blocklist", [])} + allowlist = {normalize_symbol(item) for item in settings.get("allowlist", [])} + if symbol in blocklist: + return "NO_TRADE_BLOCKED" + if allowlist and symbol not in allowlist: + return "NO_TRADE_NOT_ALLOWLISTED" + if evidence_count < int(settings["min_evidence_count"]): + return "NO_TRADE_NEEDS_EVIDENCE" + if source_count < int(settings["min_source_count"]): + return "NO_TRADE_NEEDS_SOURCE_DIVERSITY" + if confidence < float(settings["min_confidence"]): + return "NO_TRADE_LOW_CONFIDENCE" + if sentiment_score >= 0.25: + return "CONSIDER_BUY" + if sentiment_score <= -0.25: + return "CONSIDER_SELL" + return "WATCH" + + +def summarize_signal( + symbol: str, + sentiment_score: float, + confidence: float, + evidence_count: int, + source_count: int, + action: str, +) -> str: + direction = "positive" if sentiment_score > 0 else "negative" if sentiment_score < 0 else "neutral" + return ( + f"{symbol} has {direction} social sentiment from {evidence_count} evidence item(s) " + f"across {source_count} source(s). Confidence is {confidence:.0%}. Action: {action}." + ) + + +def get_config_if_available() -> SchwabConfig | None: + try: + return SchwabConfig.from_env() + except Exception: + return None +
--- /dev/null Thu Jan 01 00:00:00 1970 +0000 +++ b/schwab_trader/dashboard_server.py Sun Aug 02 16:48:18 2026 -0700 @@ -0,0 +1,331 @@ +from __future__ import annotations + +import argparse +import json +from http import HTTPStatus +from http.server import BaseHTTPRequestHandler, ThreadingHTTPServer +from typing import Any +from urllib.parse import parse_qs, urlparse + +from schwab_trader.dashboard import DashboardStore, EvidenceInput, get_config_if_available +from schwab_trader.schwab_client import SchwabError, get_account, get_account_numbers, get_accounts, load_tokens + + +INDEX_HTML = """<!doctype html> +<html lang="en"> +<head> + <meta charset="utf-8"> + <meta name="viewport" content="width=device-width, initial-scale=1"> + <title>Schwab Sentiment Dashboard</title> + <style> + :root { color-scheme: dark; font-family: Inter, system-ui, sans-serif; background: #0b1020; color: #eef2ff; } + body { margin: 0; } + header { padding: 24px; background: linear-gradient(135deg, #172554, #0f172a); border-bottom: 1px solid #334155; } + h1 { margin: 0 0 8px; font-size: 28px; } + main { display: grid; gap: 16px; grid-template-columns: repeat(auto-fit, minmax(340px, 1fr)); padding: 16px; } + section { background: #111827; border: 1px solid #334155; border-radius: 14px; padding: 16px; box-shadow: 0 10px 25px #0004; } + h2 { margin-top: 0; font-size: 18px; } + label { display: block; margin: 10px 0 4px; color: #cbd5e1; } + input, textarea, select, button { width: 100%; box-sizing: border-box; border-radius: 8px; border: 1px solid #475569; background: #020617; color: #f8fafc; padding: 10px; } + button { margin-top: 12px; background: #2563eb; border: 0; font-weight: 700; cursor: pointer; } + button.secondary { background: #334155; } + pre { overflow: auto; white-space: pre-wrap; word-break: break-word; background: #020617; padding: 12px; border-radius: 8px; } + table { width: 100%; border-collapse: collapse; font-size: 13px; } + th, td { border-bottom: 1px solid #334155; padding: 8px; text-align: left; vertical-align: top; } + .ok { color: #86efac; } + .warn { color: #fbbf24; } + .bad { color: #fca5a5; } + .span2 { grid-column: 1 / -1; } + </style> +</head> +<body> + <header> + <h1>Schwab Sentiment Dashboard</h1> + <div>This is local-first and read/paper-trade focused. There is no live-trade endpoint in this dashboard.</div> + </header> + <main> + <section> + <h2>Status</h2> + <button onclick="loadAll()">Refresh</button> + <pre id="status">Loading...</pre> + </section> + + <section> + <h2>Settings</h2> + <label>Profit target %</label><input id="profit_target_pct" type="number" step="0.1"> + <label>Stop loss %</label><input id="stop_loss_pct" type="number" step="0.1"> + <label>Max dollars per paper trade</label><input id="max_trade_dollars" type="number" step="1"> + <label>Minimum confidence</label><input id="min_confidence" type="number" step="0.01" min="0" max="1"> + <button onclick="saveSettings()">Save Settings</button> + <pre id="settingsResult"></pre> + </section> + + <section> + <h2>Add Social Evidence</h2> + <label>Source</label><input id="source" placeholder="reddit"> + <label>Symbol</label><input id="symbol" placeholder="AAPL"> + <label>URL</label><input id="url" placeholder="https://..."> + <label>Engagement</label><input id="engagement" type="number" value="0"> + <label>Text</label><textarea id="text" rows="5" placeholder="$AAPL looks bullish..."></textarea> + <button onclick="addEvidence()">Add Evidence</button> + <pre id="evidenceResult"></pre> + </section> + + <section> + <h2>Paper Trade</h2> + <label>Action</label><select id="paperAction"><option>BUY</option><option>SELL</option></select> + <label>Symbol</label><input id="paperSymbol" placeholder="AAPL"> + <label>Quantity</label><input id="paperQuantity" type="number" step="0.01" value="1"> + <label>Price</label><input id="paperPrice" type="number" step="0.01"> + <label>Reason</label><input id="paperReason" placeholder="manual paper trade"> + <button onclick="addPaperTrade()">Create Paper Trade</button> + <pre id="paperResult"></pre> + </section> + + <section class="span2"> + <h2>Signals</h2> + <div id="signals"></div> + </section> + + <section> + <h2>Recent Evidence</h2> + <div id="evidence"></div> + </section> + + <section> + <h2>Paper Trades</h2> + <div id="paperTrades"></div> + </section> + + <section class="span2"> + <h2>Audit Log</h2> + <div id="audit"></div> + </section> + </main> + <script> + async function api(path, options = {}) { + const response = await fetch(path, { + headers: {'Content-Type': 'application/json'}, + ...options + }); + const body = await response.json(); + if (!response.ok) throw new Error(body.error || response.statusText); + return body; + } + + function json(id, value) { + document.getElementById(id).textContent = JSON.stringify(value, null, 2); + } + + function table(rows, cols) { + if (!rows.length) return '<p class="warn">No data yet.</p>'; + return '<table><thead><tr>' + cols.map(c => `<th>${c}</th>`).join('') + + '</tr></thead><tbody>' + rows.map(row => '<tr>' + cols.map(c => `<td>${row[c] ?? ''}</td>`).join('') + '</tr>').join('') + '</tbody></table>'; + } + + async function loadAll() { + const [status, settings, signals, evidence, paperTrades, audit] = await Promise.all([ + api('/api/status'), api('/api/settings'), api('/api/signals'), + api('/api/evidence'), api('/api/paper-trades'), api('/api/audit') + ]); + json('status', status); + for (const key of ['profit_target_pct', 'stop_loss_pct', 'max_trade_dollars', 'min_confidence']) { + document.getElementById(key).value = settings[key]; + } + document.getElementById('signals').innerHTML = table(signals, ['symbol', 'action', 'sentiment_score', 'confidence', 'evidence_count', 'source_count', 'summary']); + document.getElementById('evidence').innerHTML = table(evidence, ['id', 'source', 'symbol', 'sentiment_score', 'engagement', 'text']); + document.getElementById('paperTrades').innerHTML = table(paperTrades, ['id', 'symbol', 'action', 'quantity', 'price', 'notional', 'status', 'reason']); + document.getElementById('audit').innerHTML = table(audit, ['id', 'event_type', 'message', 'created_at']); + } + + async function saveSettings() { + try { + const body = {}; + for (const key of ['profit_target_pct', 'stop_loss_pct', 'max_trade_dollars', 'min_confidence']) { + body[key] = Number(document.getElementById(key).value); + } + json('settingsResult', await api('/api/settings', {method: 'POST', body: JSON.stringify(body)})); + await loadAll(); + } catch (error) { json('settingsResult', {error: error.message}); } + } + + async function addEvidence() { + try { + const body = { + source: document.getElementById('source').value, + symbol: document.getElementById('symbol').value, + url: document.getElementById('url').value, + engagement: Number(document.getElementById('engagement').value), + text: document.getElementById('text').value + }; + json('evidenceResult', await api('/api/evidence', {method: 'POST', body: JSON.stringify(body)})); + await loadAll(); + } catch (error) { json('evidenceResult', {error: error.message}); } + } + + async function addPaperTrade() { + try { + const body = { + action: document.getElementById('paperAction').value, + symbol: document.getElementById('paperSymbol').value, + quantity: Number(document.getElementById('paperQuantity').value), + price: Number(document.getElementById('paperPrice').value), + reason: document.getElementById('paperReason').value + }; + json('paperResult', await api('/api/paper-trades', {method: 'POST', body: JSON.stringify(body)})); + await loadAll(); + } catch (error) { json('paperResult', {error: error.message}); } + } + + loadAll(); + </script> +</body> +</html> +""" + + +def create_handler(store: DashboardStore) -> type[BaseHTTPRequestHandler]: + class DashboardHandler(BaseHTTPRequestHandler): + server_version = "SchwabDashboard/0.1" + + def do_GET(self) -> None: + try: + parsed = urlparse(self.path) + query = parse_qs(parsed.query) + if parsed.path == "/": + self._send_html(INDEX_HTML) + elif parsed.path == "/api/status": + self._send_json(store.get_status()) + elif parsed.path == "/api/settings": + self._send_json(store.get_settings()) + elif parsed.path == "/api/evidence": + self._send_json(store.list_evidence()) + elif parsed.path == "/api/signals": + self._send_json(store.list_signals()) + elif parsed.path == "/api/paper-trades": + self._send_json(store.list_paper_trades()) + elif parsed.path == "/api/audit": + self._send_json(store.list_audit()) + elif parsed.path == "/api/schwab/account-numbers": + self._send_json(_load_schwab_account_numbers()) + elif parsed.path == "/api/schwab/accounts": + fields = "positions" if query.get("positions") == ["1"] else None + self._send_json(_load_schwab_accounts(fields)) + elif parsed.path == "/api/schwab/account": + account_hash = query.get("account_hash", [""])[0] + fields = "positions" if query.get("positions") == ["1"] else None + self._send_json(_load_schwab_account(account_hash, fields)) + else: + self._send_error(HTTPStatus.NOT_FOUND, "Not found") + except (ValueError, SchwabError) as error: + self._send_error(HTTPStatus.BAD_REQUEST, str(error)) + + def do_POST(self) -> None: + try: + parsed = urlparse(self.path) + payload = self._read_json() + if parsed.path == "/api/settings": + self._send_json(store.update_settings(payload)) + elif parsed.path == "/api/evidence": + result = store.add_evidence( + EvidenceInput( + source=str(payload.get("source", "")), + symbol=str(payload["symbol"]) if payload.get("symbol") else None, + url=str(payload["url"]) if payload.get("url") else None, + text=str(payload.get("text", "")), + engagement=float(payload.get("engagement", 0)), + raw=payload.get("raw") if isinstance(payload.get("raw"), dict) else None, + ) + ) + self._send_json(result, HTTPStatus.CREATED) + elif parsed.path == "/api/paper-trades": + self._send_json(store.add_paper_trade(payload), HTTPStatus.CREATED) + else: + self._send_error(HTTPStatus.NOT_FOUND, "Not found") + except (ValueError, KeyError, SchwabError) as error: + self._send_error(HTTPStatus.BAD_REQUEST, str(error)) + + def log_message(self, format: str, *args: Any) -> None: + print(f"[dashboard] {self.address_string()} - {format % args}") + + def _read_json(self) -> dict[str, Any]: + length = int(self.headers.get("Content-Length", "0")) + if length <= 0: + return {} + raw = self.rfile.read(length).decode("utf-8") + value = json.loads(raw) + if not isinstance(value, dict): + raise ValueError("JSON body must be an object") + return value + + def _send_html(self, body: str) -> None: + data = body.encode("utf-8") + self.send_response(HTTPStatus.OK) + self.send_header("Content-Type", "text/html; charset=utf-8") + self.send_header("Content-Length", str(len(data))) + self.end_headers() + self.wfile.write(data) + + def _send_json(self, body: Any, status: HTTPStatus = HTTPStatus.OK) -> None: + data = json.dumps(body, indent=2, sort_keys=True).encode("utf-8") + self.send_response(status) + self.send_header("Content-Type", "application/json; charset=utf-8") + self.send_header("Content-Length", str(len(data))) + self.end_headers() + self.wfile.write(data) + + def _send_error(self, status: HTTPStatus, message: str) -> None: + self._send_json({"error": message}, status) + + return DashboardHandler + + +def _load_schwab_account_numbers() -> Any: + config, access_token = _load_schwab_access_token() + return get_account_numbers(access_token).body + + +def _load_schwab_accounts(fields: str | None) -> Any: + config, access_token = _load_schwab_access_token() + return get_accounts(access_token, fields).body + + +def _load_schwab_account(account_hash: str, fields: str | None) -> Any: + if not account_hash: + raise ValueError("account_hash is required") + config, access_token = _load_schwab_access_token() + return get_account(access_token, account_hash, fields).body + + +def _load_schwab_access_token() -> tuple[Any, str]: + config = get_config_if_available() + if config is None: + raise SchwabError("Schwab environment is not configured") + tokens = load_tokens(config.token_file) + access_token = tokens.get("access_token") + if not access_token: + raise SchwabError(f"No access_token in {config.token_file}") + return config, access_token + + +def run(host: str, port: int, db_path: str | None = None) -> None: + store = DashboardStore(db_path) + server = ThreadingHTTPServer((host, port), create_handler(store)) + print(f"Schwab dashboard listening on http://{host}:{port}") + print("Live trading is disabled in this dashboard.") + server.serve_forever() + + +def main() -> int: + parser = argparse.ArgumentParser(description="Run the local Schwab sentiment dashboard") + parser.add_argument("--host", default="127.0.0.1") + parser.add_argument("--port", default=8765, type=int) + parser.add_argument("--db", help="SQLite dashboard DB path") + args = parser.parse_args() + run(args.host, args.port, args.db) + return 0 + + +if __name__ == "__main__": + raise SystemExit(main()) +
--- /dev/null Thu Jan 01 00:00:00 1970 +0000 +++ b/schwab_trader/dashboard_test.py Sun Aug 02 16:48:18 2026 -0700 @@ -0,0 +1,61 @@ +import json +import tempfile +import threading +import unittest +import urllib.request +from http.server import ThreadingHTTPServer +from pathlib import Path + +from schwab_trader.dashboard import DashboardStore, EvidenceInput, score_sentiment +from schwab_trader.dashboard_server import create_handler + + +class DashboardStoreTest(unittest.TestCase): + def test_sentiment_scoring(self): + self.assertGreater(score_sentiment("$AAPL bullish strong growth"), 0) + self.assertLess(score_sentiment("$AAPL bearish weak sell"), 0) + + def test_evidence_recomputes_signal(self): + with tempfile.TemporaryDirectory() as temp_dir: + store = DashboardStore(Path(temp_dir) / "dashboard.db") + store.update_settings({"min_evidence_count": 2, "min_source_count": 2, "min_confidence": 0.3}) + + store.add_evidence(EvidenceInput(source="reddit", symbol="AAPL", text="$AAPL bullish strong", engagement=20)) + result = store.add_evidence(EvidenceInput(source="x", symbol="AAPL", text="$AAPL growth surge", engagement=50)) + + self.assertEqual("AAPL", result["symbol"]) + signals = store.list_signals() + self.assertEqual(1, len(signals)) + self.assertEqual("CONSIDER_BUY", signals[0]["action"]) + + def test_paper_trade_rejects_above_max_trade_dollars(self): + with tempfile.TemporaryDirectory() as temp_dir: + store = DashboardStore(Path(temp_dir) / "dashboard.db") + store.update_settings({"max_trade_dollars": 100}) + + trade = store.add_paper_trade( + {"symbol": "AAPL", "action": "BUY", "quantity": 2, "price": 75, "reason": "test"} + ) + + self.assertEqual("rejected_max_trade_dollars", trade["status"]) + + def test_dashboard_http_status(self): + with tempfile.TemporaryDirectory() as temp_dir: + store = DashboardStore(Path(temp_dir) / "dashboard.db") + server = ThreadingHTTPServer(("127.0.0.1", 0), create_handler(store)) + thread = threading.Thread(target=server.serve_forever, daemon=True) + thread.start() + try: + url = f"http://127.0.0.1:{server.server_port}/api/status" + with urllib.request.urlopen(url, timeout=5) as response: + body = json.loads(response.read().decode("utf-8")) + self.assertEqual("schwab-dashboard", body["service"]) + self.assertFalse(body["live_trading_enabled"]) + finally: + server.shutdown() + thread.join(timeout=5) + + +if __name__ == "__main__": + unittest.main() +
--- /dev/null Thu Jan 01 00:00:00 1970 +0000 +++ b/schwab_trader/env.example Sun Aug 02 16:48:18 2026 -0700 @@ -0,0 +1,10 @@ +# Copy this to your shell profile or an untracked local env file. +# Never commit real Schwab credentials or OAuth tokens. + +export SCHWAB_APP_KEY="your-schwab-app-key" +export SCHWAB_APP_SECRET="your-schwab-app-secret" +export SCHWAB_REDIRECT_URI="https://127.0.0.1" + +# Optional. Defaults to ~/.config/zenbu/schwab_tokens.json. +export SCHWAB_TOKEN_FILE="$HOME/.config/zenbu/schwab_tokens.json" +
--- /dev/null Thu Jan 01 00:00:00 1970 +0000 +++ b/schwab_trader/schwab_cli.py Sun Aug 02 16:48:18 2026 -0700 @@ -0,0 +1,171 @@ +from __future__ import annotations + +import argparse +import json +import sys + +from schwab_trader.schwab_client import ( + SchwabConfig, + SchwabError, + build_authorization_url, + build_equity_order, + exchange_code_for_tokens, + get_account, + get_account_numbers, + get_accounts, + load_tokens, + place_order, + refresh_tokens, + save_tokens, +) + + +def main(argv: list[str] | None = None) -> int: + parser = argparse.ArgumentParser( + description="Safe Schwab Trader API helper. This tool never chooses trades for you.", + ) + subparsers = parser.add_subparsers(dest="command", required=True) + + auth_url = subparsers.add_parser("auth-url", help="Print the Schwab OAuth authorization URL") + auth_url.add_argument("--state", help="Optional OAuth state value") + + token = subparsers.add_parser("token", help="Exchange an authorization code/callback URL for tokens") + token.add_argument("--code", required=True, help="Authorization code or full callback URL") + + subparsers.add_parser("refresh", help="Refresh and save tokens") + + subparsers.add_parser("account-numbers", help="Print Schwab account-number/account-hash mapping") + + accounts = subparsers.add_parser("accounts", help="Print account data") + accounts.add_argument("--positions", action="store_true", help="Include positions if API permissions allow it") + + account = subparsers.add_parser("account", help="Print one account by account hash") + account.add_argument("--account-hash", required=True) + account.add_argument("--positions", action="store_true", help="Include positions if API permissions allow it") + + build_order = subparsers.add_parser("build-equity-order", help="Build and print a stock order JSON payload") + _add_order_args(build_order) + + place_equity = subparsers.add_parser( + "place-equity-order", + help="Place a user-specified stock order; defaults to dry-run output only", + ) + _add_order_args(place_equity) + place_equity.add_argument("--account-hash", required=True) + place_equity.add_argument("--live", action="store_true", help="Actually submit the order to Schwab") + place_equity.add_argument( + "--confirm-live-trade", + action="store_true", + help="Required with --live to reduce accidental orders", + ) + + args = parser.parse_args(argv) + + try: + if args.command == "auth-url": + config = SchwabConfig.from_env() + print(build_authorization_url(config.app_key, config.redirect_uri, args.state)) + return 0 + + if args.command == "token": + config = SchwabConfig.from_env() + tokens = exchange_code_for_tokens(config, args.code) + save_tokens(config.token_file, tokens) + print(f"Saved tokens to {config.token_file}") + return 0 + + if args.command == "refresh": + config = SchwabConfig.from_env() + tokens = refresh_tokens(config) + save_tokens(config.token_file, tokens) + print(f"Refreshed tokens in {config.token_file}") + return 0 + + if args.command == "account-numbers": + config = SchwabConfig.from_env() + access_token = _load_access_token(config) + _print_json(get_account_numbers(access_token).body) + return 0 + + if args.command == "accounts": + config = SchwabConfig.from_env() + access_token = _load_access_token(config) + fields = "positions" if args.positions else None + _print_json(get_accounts(access_token, fields).body) + return 0 + + if args.command == "account": + config = SchwabConfig.from_env() + access_token = _load_access_token(config) + fields = "positions" if args.positions else None + _print_json(get_account(access_token, args.account_hash, fields).body) + return 0 + + if args.command == "build-equity-order": + order = _build_order_from_args(args) + _print_json(order) + return 0 + + if args.command == "place-equity-order": + order = _build_order_from_args(args) + if not args.live: + print("DRY RUN: order was not sent. Add --live --confirm-live-trade to submit.") + _print_json(order) + return 0 + if not args.confirm_live_trade: + raise SchwabError("--live requires --confirm-live-trade") + + config = SchwabConfig.from_env() + access_token = _load_access_token(config) + response = place_order(access_token, args.account_hash, order) + print(f"Schwab order response status: {response.status}") + location = response.headers.get("Location") + if location: + print(f"Order location: {location}") + if response.body is not None: + _print_json(response.body) + return 0 + + raise SchwabError(f"Unknown command: {args.command}") + except SchwabError as error: + print(f"error: {error}", file=sys.stderr) + return 1 + + +def _add_order_args(parser: argparse.ArgumentParser) -> None: + parser.add_argument("--action", required=True, choices=["BUY", "SELL", "buy", "sell"]) + parser.add_argument("--symbol", required=True) + parser.add_argument("--quantity", required=True, type=float) + parser.add_argument("--order-type", default="MARKET", choices=["MARKET", "LIMIT", "market", "limit"]) + parser.add_argument("--price", type=float, help="Required for LIMIT orders; invalid for MARKET orders") + parser.add_argument("--duration", default="DAY") + parser.add_argument("--session", default="NORMAL") + + +def _build_order_from_args(args: argparse.Namespace) -> dict: + return build_equity_order( + action=args.action, + symbol=args.symbol, + quantity=args.quantity, + order_type=args.order_type, + price=args.price, + duration=args.duration, + session=args.session, + ) + + +def _load_access_token(config: SchwabConfig) -> str: + tokens = load_tokens(config.token_file) + access_token = tokens.get("access_token") + if not access_token: + raise SchwabError(f"No access_token in {config.token_file}") + return access_token + + +def _print_json(value: object) -> None: + print(json.dumps(value, indent=2, sort_keys=True)) + + +if __name__ == "__main__": + raise SystemExit(main()) +
--- /dev/null Thu Jan 01 00:00:00 1970 +0000 +++ b/schwab_trader/schwab_client.py Sun Aug 02 16:48:18 2026 -0700 @@ -0,0 +1,275 @@ +from __future__ import annotations + +import base64 +import json +import os +import stat +import time +import urllib.error +import urllib.parse +import urllib.request +from dataclasses import dataclass +from pathlib import Path +from typing import Any + + +AUTH_URL = "https://api.schwabapi.com/v1/oauth/authorize" +TOKEN_URL = "https://api.schwabapi.com/v1/oauth/token" +TRADER_BASE_URL = "https://api.schwabapi.com/trader/v1" +DEFAULT_TOKEN_FILE = "~/.config/zenbu/schwab_tokens.json" + + +class SchwabError(RuntimeError): + pass + + +@dataclass(frozen=True) +class SchwabConfig: + app_key: str + app_secret: str + redirect_uri: str + token_file: Path + + @classmethod + def from_env(cls) -> "SchwabConfig": + app_key = os.environ.get("SCHWAB_APP_KEY", "").strip() + app_secret = os.environ.get("SCHWAB_APP_SECRET", "").strip() + redirect_uri = os.environ.get("SCHWAB_REDIRECT_URI", "").strip() + token_file = Path(os.environ.get("SCHWAB_TOKEN_FILE", DEFAULT_TOKEN_FILE)).expanduser() + + missing = [ + name + for name, value in ( + ("SCHWAB_APP_KEY", app_key), + ("SCHWAB_APP_SECRET", app_secret), + ("SCHWAB_REDIRECT_URI", redirect_uri), + ) + if not value + ] + if missing: + raise SchwabError("Missing required environment variables: " + ", ".join(missing)) + + return cls( + app_key=app_key, + app_secret=app_secret, + redirect_uri=redirect_uri, + token_file=token_file, + ) + + +@dataclass(frozen=True) +class ApiResponse: + status: int + headers: dict[str, str] + body: Any + raw_body: str + + +def build_authorization_url(app_key: str, redirect_uri: str, state: str | None = None) -> str: + params = { + "response_type": "code", + "client_id": app_key, + "redirect_uri": redirect_uri, + } + if state: + params["state"] = state + return AUTH_URL + "?" + urllib.parse.urlencode(params) + + +def extract_authorization_code(code_or_url: str) -> str: + value = code_or_url.strip() + if not value: + raise SchwabError("Authorization code is empty") + + parsed = urllib.parse.urlparse(value) + if parsed.scheme and parsed.netloc: + query = urllib.parse.parse_qs(parsed.query) + codes = query.get("code") + if not codes or not codes[0]: + raise SchwabError("No code= parameter found in callback URL") + return codes[0] + + return urllib.parse.unquote(value) + + +def exchange_code_for_tokens(config: SchwabConfig, code_or_url: str) -> dict[str, Any]: + code = extract_authorization_code(code_or_url) + return _token_request( + config, + { + "grant_type": "authorization_code", + "code": code, + "redirect_uri": config.redirect_uri, + }, + ) + + +def refresh_tokens(config: SchwabConfig, refresh_token: str | None = None) -> dict[str, Any]: + token_value = refresh_token + if token_value is None: + existing = load_tokens(config.token_file) + token_value = existing.get("refresh_token") + if not token_value: + raise SchwabError("No refresh token available") + + return _token_request( + config, + { + "grant_type": "refresh_token", + "refresh_token": token_value, + }, + ) + + +def save_tokens(path: Path, tokens: dict[str, Any]) -> None: + payload = dict(tokens) + payload["saved_at"] = int(time.time()) + path.parent.mkdir(parents=True, exist_ok=True) + path.write_text(json.dumps(payload, indent=2, sort_keys=True) + "\n", encoding="utf-8") + path.chmod(stat.S_IRUSR | stat.S_IWUSR) + + +def load_tokens(path: Path) -> dict[str, Any]: + if not path.exists(): + raise SchwabError(f"Token file does not exist: {path}") + return json.loads(path.read_text(encoding="utf-8")) + + +def get_account_numbers(access_token: str) -> ApiResponse: + return _api_request("GET", "/accounts/accountNumbers", access_token) + + +def get_accounts(access_token: str, fields: str | None = None) -> ApiResponse: + query = "" + if fields: + query = "?" + urllib.parse.urlencode({"fields": fields}) + return _api_request("GET", "/accounts" + query, access_token) + + +def get_account(access_token: str, account_hash: str, fields: str | None = None) -> ApiResponse: + query = "" + if fields: + query = "?" + urllib.parse.urlencode({"fields": fields}) + return _api_request("GET", f"/accounts/{urllib.parse.quote(account_hash)}{query}", access_token) + + +def build_equity_order( + action: str, + symbol: str, + quantity: float, + order_type: str = "MARKET", + price: float | None = None, + duration: str = "DAY", + session: str = "NORMAL", +) -> dict[str, Any]: + normalized_action = action.upper() + normalized_symbol = symbol.upper() + normalized_order_type = order_type.upper() + normalized_duration = duration.upper() + normalized_session = session.upper() + + if normalized_action not in {"BUY", "SELL"}: + raise SchwabError("action must be BUY or SELL") + if not normalized_symbol: + raise SchwabError("symbol is required") + if quantity <= 0: + raise SchwabError("quantity must be greater than zero") + if normalized_order_type not in {"MARKET", "LIMIT"}: + raise SchwabError("order_type must be MARKET or LIMIT") + if normalized_order_type == "LIMIT" and price is None: + raise SchwabError("LIMIT orders require --price") + if normalized_order_type == "MARKET" and price is not None: + raise SchwabError("MARKET orders cannot include --price") + + order: dict[str, Any] = { + "orderType": normalized_order_type, + "session": normalized_session, + "duration": normalized_duration, + "orderStrategyType": "SINGLE", + "orderLegCollection": [ + { + "instruction": normalized_action, + "quantity": quantity, + "instrument": { + "symbol": normalized_symbol, + "assetType": "EQUITY", + }, + } + ], + } + if price is not None: + order["price"] = f"{price:.2f}" + return order + + +def place_order(access_token: str, account_hash: str, order: dict[str, Any]) -> ApiResponse: + path = f"/accounts/{urllib.parse.quote(account_hash)}/orders" + return _api_request("POST", path, access_token, order) + + +def _token_request(config: SchwabConfig, form: dict[str, str]) -> dict[str, Any]: + credentials = f"{config.app_key}:{config.app_secret}".encode("utf-8") + headers = { + "Authorization": "Basic " + base64.b64encode(credentials).decode("ascii"), + "Content-Type": "application/x-www-form-urlencoded", + "Accept": "application/json", + } + data = urllib.parse.urlencode(form).encode("utf-8") + request = urllib.request.Request(TOKEN_URL, data=data, headers=headers, method="POST") + response = _open_request(request) + if not isinstance(response.body, dict): + raise SchwabError("Token endpoint did not return a JSON object") + return response.body + + +def _api_request( + method: str, + path: str, + access_token: str, + body: dict[str, Any] | None = None, +) -> ApiResponse: + headers = { + "Authorization": f"Bearer {access_token}", + "Accept": "application/json", + } + data = None + if body is not None: + data = json.dumps(body).encode("utf-8") + headers["Content-Type"] = "application/json" + + request = urllib.request.Request( + TRADER_BASE_URL + path, + data=data, + headers=headers, + method=method, + ) + return _open_request(request) + + +def _open_request(request: urllib.request.Request) -> ApiResponse: + try: + with urllib.request.urlopen(request, timeout=30) as response: + raw = response.read().decode("utf-8") + return ApiResponse( + status=response.status, + headers=dict(response.headers.items()), + body=_parse_json_or_text(raw), + raw_body=raw, + ) + except urllib.error.HTTPError as error: + raw = error.read().decode("utf-8", errors="replace") + raise SchwabError( + f"Schwab API request failed with HTTP {error.code}: {raw or error.reason}" + ) from error + except urllib.error.URLError as error: + raise SchwabError(f"Schwab API request failed: {error.reason}") from error + + +def _parse_json_or_text(raw: str) -> Any: + if not raw: + return None + try: + return json.loads(raw) + except json.JSONDecodeError: + return raw +
--- /dev/null Thu Jan 01 00:00:00 1970 +0000 +++ b/schwab_trader/schwab_client_test.py Sun Aug 02 16:48:18 2026 -0700 @@ -0,0 +1,67 @@ +import os +import tempfile +import unittest +from pathlib import Path +from urllib.parse import parse_qs, urlparse + +from schwab_trader.schwab_client import ( + SchwabConfig, + SchwabError, + build_authorization_url, + build_equity_order, + extract_authorization_code, + save_tokens, +) + + +class SchwabClientTest(unittest.TestCase): + def test_authorization_url_contains_oauth_params(self): + url = build_authorization_url("app-key", "https://127.0.0.1/callback", "state-1") + parsed = urlparse(url) + params = parse_qs(parsed.query) + + self.assertEqual("https", parsed.scheme) + self.assertEqual("api.schwabapi.com", parsed.netloc) + self.assertEqual(["code"], params["response_type"]) + self.assertEqual(["app-key"], params["client_id"]) + self.assertEqual(["https://127.0.0.1/callback"], params["redirect_uri"]) + self.assertEqual(["state-1"], params["state"]) + + def test_extract_authorization_code_from_callback_url(self): + code = extract_authorization_code("https://127.0.0.1/callback?code=abc%40123&state=x") + self.assertEqual("abc@123", code) + + def test_build_market_equity_order(self): + order = build_equity_order("buy", "aapl", 1) + + self.assertEqual("MARKET", order["orderType"]) + self.assertEqual("BUY", order["orderLegCollection"][0]["instruction"]) + self.assertEqual("AAPL", order["orderLegCollection"][0]["instrument"]["symbol"]) + self.assertNotIn("price", order) + + def test_limit_order_requires_price(self): + with self.assertRaises(SchwabError): + build_equity_order("SELL", "MSFT", 2, order_type="LIMIT") + + def test_save_tokens_uses_owner_only_permissions(self): + with tempfile.TemporaryDirectory() as temp_dir: + token_file = Path(temp_dir) / "tokens.json" + save_tokens(token_file, {"access_token": "x"}) + + self.assertEqual(0o600, token_file.stat().st_mode & 0o777) + + def test_config_from_env_requires_credentials(self): + old_env = os.environ.copy() + try: + for key in ("SCHWAB_APP_KEY", "SCHWAB_APP_SECRET", "SCHWAB_REDIRECT_URI"): + os.environ.pop(key, None) + with self.assertRaises(SchwabError): + SchwabConfig.from_env() + finally: + os.environ.clear() + os.environ.update(old_env) + + +if __name__ == "__main__": + unittest.main() +