Skip to content

Latest commit

 

History

History
458 lines (414 loc) · 28.6 KB

File metadata and controls

458 lines (414 loc) · 28.6 KB

The web layer

web_interface/ contains the Flask dashboard, its HTTP API, and the background workers. The full endpoint list is in routes.md (generated by scripts/gen_route_inventory.py).

Package Holds
fyp_data_hub.py, static_assets.py, seo.py the app factory and its app-wide helpers
routes/ the Flask blueprints (below) — the HTTP surface only
services/ business logic shared by routes and workers
auth/ accounts and access control (below)
tasks/ the background-task runtime: worker registry, launch, status, logs, failure ledger, the Cloud Tasks runtime (Background workers)
workers/ one run_<name>.py module per background worker
integrations/ outbound email (mail_utils) and Slack (slack_service)

Layering: routes and workers import services; services never import a route module or the Flask-facing auth modules (auth.security, auth.permissions) — they may use the account store, auth.accounts. Workers never import a route module.

App structure

fyp_data_hub.py is an app factory (create_app). It registers 13 blueprints from web_interface/routes/ — 12 unconditionally, plus the CSRF-exempt internal_bp only on the task-runner service or outside Cloud Run (see create_app in fyp_data_hub.py):

Blueprint file Serves
auth_routes/ (package) login, signup, email verification, logout (login); admin users, roles and site settings (admin_users, admin_roles, admin_site); a user's own settings and profile (user)
public_routes.py the public (unauthenticated) mini-site: landing, about, participate + the /participate/start wizard, data-donation, thehub, terms, ethics, faq, robots.txt, sitemap.xml, and legacy-URL 301s (/guide → /thehub, old Wix paths)
my_collections_routes.py my_collections_bp — the participant self-service API under /api/my/collections/*: list, upload sources, upload, pending personality/delete, withdraw, restore, process, per-collection and combined personality
api_explorer_routes/ (package) studies + Explore tab API (explore); admin System Information, System Health and the ops report (system)
api_viewer_routes.py Video Analysis tab + media streaming
api_timelines_routes.py Timelines tab
api_correlations_routes.py Correlations tab
api_semantic_space_routes.py Semantic Space tab (embedding map)
api_sessions_routes.py Sessions tab (session index + binge episodes + low-entropy sequences); its data and statistics are services/sessions_data.py and services/sessions_stats.py
management/ (package) Data Pipeline + admin: studies, collections, enrichment queues, contracts, data contracts, A/B evaluation, schema, ingestion — split into per-domain submodules all registering on the same blueprint
human_eval_routes.py human annotation input (coding, votes, invitations)
process_routes.py background-process control + the CSRF-exempt internal_bp that receives Cloud Tasks pushes at /internal/run-task/<name> (the runtime behind it is tasks/runtime.py)

The services/ package holds study data + cache (note the double-checked locking in StudyCache), timelines, analysis data, correlations, stats, per-user variables, worker status, preview cache, system health, the per-study methods/provenance note builder methods_note.py, the participant surface (my_collections_service.py, participant_studies.py, participant_enrichment.py), collection_coverage.py (the shared scraped/annotated-coverage arithmetic, so the participant and admin coverage figures cannot drift), collection_enrichment.py (the plan ledger and slice cutter behind the automatic enrichment loop), enrichment_journal.py (the durable enrichment history), refresh_pipeline.py (the refresh-run dependency registry, run record and planner), downstream_refresh.py (dispatches the downstream refresh pipeline from a consolidation's impact), collection_deletion.py, the daily ops_report.py, and the explorer backend, collection-trajectory overlay, collection ↔ account links, admin settings, admin notes, activity log and citation modules.

Auth & permissions

Flask-Login over a JSON-file user store, all in web_interface/auth/: accounts.py is the store (UserManager, RoleManager, the user_manager singleton; user records are JSON files under {local_data}/users/), security.py wires it into Flask-Login, and permissions.py holds the tab/sub-page permission catalog and the route guards. Route guards: @permission_required("<key>") with a key from PERMISSION_CATALOG (admins pass every check; it already sends an unauthenticated request to the login flow, so it is never stacked with @login_required), @login_required alone for any-signed-in-user routes, and @admin_required for the few admin-only endpoints that have no catalog key. CSRF is globally enabled (Flask-WTF); only the OIDC-authenticated internal task blueprint is exempt. WTF_CSRF_TIME_LIMIT = None is deliberate (long-open research sessions).

Three built-in roles are seeded at boot: admin ("*"), viewer (the analysis tabs + personal My-stuff pages), and the read-only student teaching role (permissions.STUDENT_PERMISSIONS: the viewer set minus Semantic Space, Sessions and feature.annotation_votes — the key gating both vote endpoints; the boot migration grant-alls it to existing roles but skip-lists student). Per-study sharing is the study definition's USER_ACCESS list (role names / usernames / 'all'); an empty or missing list means shared with nobody on every surface, so sharing is always an explicit grant. A boot-time migration (fyp.analysis.studies.migrate_user_access_defaults, serving processes only) backfills explicit grants into studies from before that rule (see decision 0007).

data.sensitive_activity ("Sensitive activity — read comment text" on the User Roles page) gates the text of comments donors wrote. It is in no default set and no boot-time grant, so every non-admin role starts without it. A role without it still sees that a comment was made (the bare comment engagement token, its counts and the engagement filter), never what was said: study_data.get_explorer_data / get_explorer_rows strip comment:<text> tokens from extra_data unless the caller passes hide_comment_text=False (the per-user routes pass not can_read_sensitive_activity(current_user)), and global search runs on the stripped column. Summaries never carry the raw cells for anyone: extra_data stats are per-token counts, it is never a Timelines series, and the My Collections favourite emoji reaches only the donor and key holders.

Participants who own donated collections additionally get an auto-managed study pair — __me__{username} ("Just Me", their own collections, materialised) and __me_plus__{username} ("Everyone & Me", composed at read time from the site default study plus the Just Me dataset; never materialised). The pair is created/updated by services/participant_studies.py from the ownership links in collections_tags.json. Provisioning is lazy: an owner's pair is first created on login (ensure_on_login, run on every successful login), so donation-linked accounts that never log in get no studies; existing pairs are reconciled by every sync path regardless — the other lifecycle callers are run_ingest_refresh after a donation registers, the admin collection re-link route (after collection_accounts.set_collection_owner), run_collection_delete, and account deletion in the auth routes. Defs carry SYSTEM/OWNER/DISPLAY_NAME markers and USER_ACCESS = [username] (__me_plus__ defs additionally carry a COMPOSE marker — services/study_data.py assembles the composed frame at read time). The pair is skipped by the boot migration and the all-studies refresh sweeps, is excluded from the default-study picker, and is refused by the study save/rename endpoints (delete is admin-only cleanup). POST /api/user/settings accepts only the whitelisted USER_SETTINGS_KEYS.

A user record also carries a profile block (auth.PROFILE_FIELDS: full name, age, postcode, country, occupation, TikTok handle, consent to contact; validated by auth.validate_profile), an account_kind (member — signed up or admin-created — or participant — created from donation data, initially with no password and so unable to log in until an admin sets one), a placeholder flag for the fake p-N@<domain> addresses minted for participants who left demographics but no email, an origin, and a terms_accepted_at timestamp set at signup when the terms checkbox is ticked. Signing up with the email of a passwordless participant account claims it (UserManager.claim_participant_account).

Collections ↔ accounts. A collection belongs to at most one account. The link is the user_id key of the collection's entry in recoded/collections_tags.json — absent = undecided (AIO ingest may auto-link), null = explicitly unassigned (never re-linked), a username = linked. web_interface/services/collection_accounts.py owns the format (set_collection_owner, unlink_user, orphan_placeholder_accounts), the AIO donor-data → account move (link_aio_collections, run by the ingest worker after save_processed and by the one-off migration migrate_existing_collections), and the rule that the AIO demographic fields (fyp.ingest.donations.AIO_DEMOGRAPHIC_FIELDS) are stripped by every writer of the collections metadata parquet. Pickers: GET /api/manage/accounts; the link is set via the upload route (user_id form field) and POST /api/manage/collection/save_annotation.

Background workers

Background jobs run the same code in two modes. On Cloud Run, eligible processes run as Google Cloud Tasks dispatched to the fyp-task-runner service; locally (and for anything not Cloud-Task-eligible) they run as subprocesses. The switch is automatic, on the K_SERVICE environment variable.

Each workers/run_<name>.py module is dual-mode:

  • a run_<name>(reporter, task_args) function invoked by Cloud Tasks via tasks/runtime.TASK_FUNCTIONS (reporter = GCSStatusReporter, which writes progress and data to GCS status files, task_status/*.json, with a background heartbeat thread every 30 s), and
  • a __main__ block for local subprocess mode, run as python -m web_interface.workers.run_<name> (reporter = LocalStatusReporter, which prints ::PROGRESS:: / ::DATA:: lines that process_manager.py parses from stdout — do not break this contract; see CONTRIBUTING.md invariant #2).

Everything the two modes need to know about a worker — its name, module, Cloud Tasks dispatch deadline, whether the queue may retry it, and which launch surfaces offer it — is declared once, in WORKERS in web_interface/tasks/worker_registry.py; process_manager.CLOUD_TASK_ELIGIBLE, tasks/runtime.TASK_FUNCTIONS / QUEUE_RETRY_SAFE and the local launch command are derived from it (guard: tests/unit/test_worker_registry.py). A new worker is one run_<name>.py module plus one WORKERS entry.

Dispatch. process_manager.start_process() picks the mode: on Cloud Run it dispatches through dispatch_cloud_task(), locally it spawns a subprocess. Cloud Tasks deliver to the CSRF-exempt internal_bp blueprint at /internal/run-task/<name>, which authenticates the request and hands it to tasks/runtime.run_task_with_stats() for execution, stats and chaining.

  • Subprocess mode pins the child's project root (process_manager.worker_env()): the spawned worker gets FYP_CONFIG_PATH set to the config TOML the server itself loaded, plus PROJECT_ROOT prepended to PYTHONPATH. Without both, the child rediscovers its own root — fyp.core.paths walks up from the working directory for __proj__.py, and import fyp can be answered by the venv's editable install pointing at a different checkout — and can load another config.local.toml, and so another data store, than its parent (see decision 0012; guard: tests/unit/test_worker_spawn_env.py).
  • Explicit dispatch deadlines. Cloud Tasks' default HTTP deadline is shorter than a heavy batch link, and a timed-out attempt keeps running while the retry starts a concurrent duplicate chain. Every self-chaining refresh therefore carries an explicit 1800 s deadline, declared once per worker in WORKERS and read through worker_registry.deadline_for() by every dispatch — the initial one from process_manager and each link a worker chains (see decision 0008; tests/unit/test_dispatch_deadlines.py pins the table).
  • Single-flight leases. embeddings_refresh claims a CAS-guarded lease file so a Cloud Tasks redelivery can never run two appenders against the embedding shard store at once; the store readers additionally dedupe on item id, last occurrence wins (see decision 0009).

Self-chaining. Long-running workers (the annotator, the scrapers, the heavy refreshes) process one batch per Cloud Task and return {"chain": True, "next_task_args": ...} to dispatch the next; each link inherits the GCS status via reporter.resume(). A task whose GCS heartbeat is older than 600 s is treated as dead by the UI and by start_process().

Retries. The Cloud Tasks queue (fyp-background-tasks) allows up to four attempts with backoff, but retry is app-controlled: only the idempotent refreshes in tasks/runtime.QUEUE_RETRY_SAFE answer a failure with 503 (and are retried); every other task answers 200 and its failure is terminal. All failures land in the task-failures ledger (cache/task_failures.json, web_interface/tasks/task_failures.py) — the dead-letter record, surfaced on Admin → System Information. Queue setup: scripts/configure_task_queue.sh (see DEVELOPING.md).

Status and logs. The reported last_run_duration spans the whole run of a self-chaining task, not just its final link: tasks/runtime._chain_run_start() measures from the first link's start (carried through task_args), and api_status forwards the queued/failed GCS states so the UI's status lights (one green/blue/amber/red vocabulary, see Process UI below) reflect them. Every run also lands in a durable per-process log (run_logs.py, not a worker): a ring of the last 10 runs per process (proc_logs/<status_key>.json in the "cache" location), each run timestamped line by line (once, in append()) and opened with a Started by <user> banner. Both execution modes and both Cloud Run services write the same document (CAS writes + a per-key flusher thread), so logs survive restarts and are visible from either service and to every admin; log bookkeeping never raises into the task. GET /api/logs/<name> serves them to the log modal, and POST /api/logs/clear/<name> empties a process's ring.

Cross-service state. Both Cloud Run services share process_stats.json on GCS; always call load_process_stats() before reading or writing it, so one service never clobbers the other's data.

Two workers are notable for how they deviate:

  • enrichment_supervisor — the automatic per-collection enrichment loop (run_enrichment_supervisor.py + services/collection_enrichment.py; the loop's rules are in pipeline.md). It is deliberately a conductor, not an executor: the scrape queues, the annotation queue, the queue workers and the consolidation pipeline are global singletons, so each short tick starts at most one of them and returns rather than doing the work itself. It never self-chains, which keeps it clear of the dispatch-deadline trap; progress comes from three idempotent triggers (a terminal worker completion, the end of a consolidation, and an hourly Cloud Scheduler heartbeat). Every tick re-reads the world and defers while any enrichment or pipeline step is running. What it and the workers do is written to the enrichment history (services/enrichment_journal.py, a bounded ring in cache/enrichment_journal.json), shown on Dataset Assembly and, per collection, in the Edit Collections panel.
  • ops_report — the daily operational health report (run_ops_report.py
    • services/ops_report.py): checks across the whole system, an AI-written assessment, and an emailed copy; artifacts land under cache/ops_report/. It is the one worker that is deliberately not queue-retry-safe — a queue retry would re-send the email — so a failed run goes straight to the task-failures ledger instead of being retried.

Frontend

No-build-step SPA: templates/index.html is the shell, templates/tabs/ holds per-tab content, static/main.js is the tab-navigation controller and each tab has its own JS file. Everything is vanilla JS + fetch(); scripts are plain <script> tags. Templates reference every script and stylesheet as {{ asset_url('main.js') }} (web_interface/static_assets.py), which appends a hash of the file's content — a changed file gets a new URL on its own, so there is no version to bump. Because of that, a response whose hash matches the file served is sent as Cache-Control: public, max-age=31536000, immutable, and browsers reuse it without a revalidation request (a mismatched hash, as during a deploy, keeps the default revalidation). Shared helpers (escapeHtml, showToast) live once in static/js/core/dom_utils.js, and the JSON fetch helpers (getJSON, postJSON: parsed body, or an Error carrying the server's message, .status and .body) in static/js/core/api.js (app pages only), both loaded in base.html <head>: all scripts share one global scope, and tests/unit/test_js_global_collisions.py fails if two files define the same top-level name.

Notable non-tab scripts in static/js/: donation_review.js (browser-side donation-export parsing and pruning — hard invariant: it makes no network requests, and the pruned file is rebuilt from kept rows only), hub_tour.js (the guided tour), admin_ops_report.js (Admin → System ops report pane), and my_collections.js, which exports window.mycRenderPersonality(); js/data_management/edit_collections.js reuses it in the Edit Collections modal so the participant and admin persona views cannot drift.

Styling is entirely token-driven (static/style.css): semantic CSS custom properties, a 7-step type scale, utility classes, and both dark (:root) and light ([data-theme="light"]) themes. Never hardcode colors/fonts/sizes in templates or JS — see DEVELOPING.md for the full rules.

Filter dropdowns for categorical/list variables show the top-200 most-frequent values (single-occurrence values are dropped entirely at metadata-build time — see explorer_backend.get_metadata). Any capped dropdown gets a per-variable search box (static/filter_value_search.js, shared by Explore and Video Analysis) backed by GET /api/explore/values/search, which matches against all values of that one column — singletons included — via a per-(study, column) counts cache in services/study_data.py (search_column_value_counts; cold path is a selective single-column parquet read that never touches StudyCache). The global free-text search is a separate mechanism: it scans every searchable column (explorer_backend.search_columns derives the set, and the filter/ids endpoints project the frame to it instead of loading full width) with one string cast per column per request.

Process UI (Data Pipeline → Dataset Assembly)

The worker cards live in the templates/tabs/dm/ partials (refresh.html and friends); a card is inert until main.js's poll loop calls setStatus() for its process name — new workers (e.g. sessions_refresh) must be wired there. Every card carries a card_info() ⓘ tooltip (one Jinja macro; the tooltip text doubles as the accessible name). Status lights use the unified --status-* tokens in style.css: green = running, blue = idle/stopped, amber = warn, red = error/failed.

Refresh runs

Starting any of these cards starts a refresh run: the steps that read what it writes follow it automatically. The graph and the rules live in one registry, services/refresh_pipeline.py, together with the predicates that decide, at each completion, whether the next step has anything to do. A step is dispatched only when an upstream step reports that something actually changed (embeddings written, videos that moved niche, study datasets rebuilt); the rule is to prune only on a positive statement of no change, so a missing signal always runs. The run is recorded in process_stats["refresh_pipeline"] and drawn as a wall-clock Gantt above the cards, with a header naming its origin, and each skipped step stating why. Only one run happens at a time: the other cards grey out while one is in flight and /api/start refuses with 423.

The Consolidate card is one origin among several. Its start dialog carries "Refresh caches afterwards" plus "Force full rebuild"; its own summary line still reports the outcome of a run it started. Scope: a run is planned for exactly what the consolidation touched — its affected studies and collections are the scope of recode, metadata, correlations and timelines, and that scope wins even when the semantic map moved videos between niches, because a warm-started rebuild moves a couple of percent of the corpus on almost every run and widening would turn every run into a full refresh (see decision 0014). The only unscoped run is one started from the Semantic Map card, which never consolidated and so has no impact to scope by. Each step records the scope it was actually dispatched with and why; hovering its bar shows it, and a pending map-dependent step states the impact's number as a floor.

Unticking "Refresh caches afterwards" means it. The consolidation still computes and reports its impact, and that impact stays on the card offering "Refresh All Affected" until the operator presses it. It is never spent by the enrichment supervisor: the deferred-impact ledger records whose debt it is (from_plan), the supervisor's own consolidations are tagged plan_deferred, and its finalize spends only those — see the enrichment loop in pipeline.md.

Liveness. The hub's abandoned-run sweep closes a run only when nothing is running, nothing has completed, and no dispatched step is still awaiting delivery, for ABANDONED_RUN_SECONDS (600). Activity is read from the workers' own status files, not from the hub's lazily-loaded process_stats, and a spine dispatch is stamped queued exactly as fork leaves are, so a task the queue holds back (a 429, redelivered on backoff) keeps its run alive for the queue's full retry window (QUEUED_DELIVERY_GRACE_SECONDS = the queue's maxRetryDuration, 3600 s). The Dataset Assembly page re-reads server state every 30 s while open, so a run the server starts on its own is noticed. sessions_refresh is also chained after every study save, which is a private chain rather than a run and stays out of the chart.

Study modal and per-user preferences

The edit-study modal's footer groups Delete | Rename | Duplicate: delete and rename have their own endpoints (POST /api/manage/studies/delete / .../rename — rename moves the study's cached artifacts in place, no rebuild), while duplicate is client-side (js/data_management/studies.js duplicateStudy() saves a copy under a suggested unique name). The study date window has per-end controls — date inputs, −/+ day steppers, and draggable chart edge handles over the daily-activities chart (POST /api/manage/studies/daily_activities serves the full-span series so the user can pick a window).

Per-user variable preferences. Each user can include, exclude and reorder variables per surface (filter / viz / detail panel / timeline) in My stuff → Preferences → Variable customizations: one dual-list dialog with a tab per surface (VariablePrefs.openCustomizer in static/js/variable_prefs.js, fed by the study-independent GET /api/user/variable-catalog). Reordering is within a section; section order stays fixed. The choices are stored as deltas in user.settings.variable_prefs (include / exclude, plus an optional order) and composed as (global ∪ include) − exclude, arranged by the surface's default order and then the user's own — client-side by VariablePrefs.effectiveFor, and server-side (compose_effective_variables) for the Timelines tab and the Explore filter-stats endpoint (/api/explore/filter computes distribution stats only for the user's effective viz set). Orders are applied by one rule on both sides (apply_section_order): within each section, the listed variables are re-sorted in place and unlisted ones keep their slot. The default order per surface (default_order in the metadata) is the computed order — section, then categorical before numerical, then A–Z — re-sorted by the admins' arrangement from Admin → Variable Visibility → Arrange default layout, which opens the same dialog in admin mode. Variables in the backstage section and those with the skip role are left out of every list (HIDDEN_SECTIONS in services/study_data.py). The same free-form settings merge (POST /api/user/settings) backs other per-user UI state, e.g. the Home tab's dismissible "Getting started" panel (getting_started_dismissed) — the permission-keyed first-run surface for newly invited users, which also links into the public /thehub page (/guide is now a 301 to it).

Conventions & known warts

  • Every error response carries an "error" key with a message for the user, with an appropriate 4xx. Some endpoints (mostly the 409 "already running" refusals of process starts) also send the older "status": "error" and "message" keys, which the JS still reads; a guard test (tests/unit/test_api_error_hygiene.py) keeps "error" on every one of them.
  • An unexpected exception never reaches the client as text. A route lets it escape (no broad except Exception returning str(e)), and the app-wide handler in fyp_data_hub.py logs the traceback under a short reference id and answers /api/* with {"error": "Internal error (ref <id>)"}, 500; search the logs for the id. A handler that must keep its own response shape calls routes/_errors.py's log_unexpected. Narrow handlers for our own exception types (a ValueError raised for bad input, a conflict error) may show their message. The same guard test enforces this.
  • Success envelopes are historically inconsistent (bare payload, {"success": true}, {"ok": true}, {"status": "success"}) — match the file you are editing; a unification is planned but is a breaking change for the JS.
  • The big historical files have been split: management routes live in routes/management/ (per-domain submodules), backend logic in services/, and the former inline template scripts in static/js/ (admin_tab.js, my_stuff_tab.js, ...). Add new admin functionality in those locations.
  • The Data Pipeline script is split by feature under static/js/data_management/ (core, studies, enrichment, ingestion, edit_collections, auto_enrichment, collection_actions). The files share one global scope and load in the order templates/index.html lists them: a top-level statement may only call functions from its own or an earlier file. Users with the Data Pipeline tab get all of them; a My Studies viewer without it gets only core and studies (the read-only study list and modal), so those two may call into the other files only on paths such a viewer never reaches — tests/unit/test_data_management_scripts.py lists the allowed calls. The worker cards, status poll and log modal are static/js/worker_control.js, loaded right after main.js.
  • A .modal dialog is built with the modal() macro in templates/_macros.html ({% call modal(id, close, ...) %}...{% endcall %}): backdrop, box and a real close <button class="modal-close">, which main.js clicks when Escape is pressed on the topmost open modal. The contract editor is the deliberate exception (no close on Escape or the backdrop, it holds unsaved work). Messages and confirmations use showAppAlert / showAppConfirm, never the browser's alert/confirm; tests/unit/test_modal_markup.py and test_js_no_native_dialogs.py enforce both.
  • Still-deferred frontend work: removing inline onclick= handlers, and hex-color/token cleanup. The inline handlers pin functions to window, so a handler's function cannot be moved into a module scope or renamed without updating every template that calls it.