From dc449be4e1eab59b73b2a9bf4d759745a7532345 Mon Sep 17 00:00:00 2001 From: Sijie Guo Date: Sat, 3 Oct 2026 01:31:10 -0700 Subject: [PATCH] Run the Cloud course in an instance per participant The Cloud course no longer starts from a team card. Each participant gets an instance in the hackathon organization and a service account from the organizers, creates a Kafka cluster, an agent workspace and a SQL workspace in it, and fills in .env with snctl lookups in Lab 0. - Lab 0 reads each address with snctl, creates the login topic and loads it. The seeders accept the cloud stack and name the course in their hints. - Lab 2 selects from the source named after the topic, security.login_events, and offers psql through the SQL workspace gateway. - Lab 4 says the SQL workspace needs read-write MCP access, and its deny example matches what the scripts print. - Troubleshooting covers the snctl context, the SQL tools failing with connection_unavailable, read-only MCP access, and a slow approval. - The doctor and the setup hints point at the instance, not a team card. - What was run: Labs 1, 3 and 4 with every check on all three paths against a test instance on StreamNative Cloud. --- .env.cloud.example | 39 ++++--- README.md | 12 +-- cli/env.sh | 4 +- cli/tests/run.sh | 2 +- docs/before-you-arrive.md | 51 +++++++-- docs/tutor.md | 2 +- labs/README.md | 10 +- labs/cloud/00-set-up.md | 160 +++++++++++++++++++++++------ labs/cloud/02-streaming-sql.md | 57 +++++----- labs/cloud/03-live-context.md | 2 +- labs/cloud/04-act-with-approval.md | 12 ++- labs/cloud/README.md | 48 +++++---- labs/cloud/troubleshooting.md | 24 +++-- local/tests/run.sh | 3 +- local/write-env.sh | 4 +- python/common.py | 10 +- python/doctor.py | 14 ++- python/seed.py | 11 +- python/tests/test_doctor.py | 49 ++++++++- python/tests/test_seed.py | 38 ++++++- python/tests/test_stack.py | 10 ++ skills/data-agent-tutor/SKILL.md | 2 +- sql/cloud/01_explore.sql | 13 +-- sql/cloud/02_login_failures.sql | 9 +- sql/cloud/03_flagged_accounts.sql | 2 +- sql/cloud/99_reset.sql | 2 +- typescript/src/common.ts | 10 +- typescript/src/doctor.ts | 39 ++++--- typescript/src/seed.ts | 39 ++++--- typescript/test/doctor.test.ts | 29 +++++- typescript/test/seed.test.ts | 49 ++++++++- typescript/test/stack.test.ts | 9 ++ 32 files changed, 562 insertions(+), 203 deletions(-) diff --git a/.env.cloud.example b/.env.cloud.example index d2d8ee8..b715036 100644 --- a/.env.cloud.example +++ b/.env.cloud.example @@ -1,29 +1,35 @@ -# The Cloud course: copy this file to .env (in the repo root) and paste the -# values from your team card. .env is git-ignored: never commit it. +# The Cloud course: copy this file to .env (in the repo root) and fill it in +# from your instance on StreamNative Cloud. .env is git-ignored: never commit it. # # cp .env.cloud.example .env # +# labs/cloud/00-set-up.md, step 2, has the snctl command for each address. # (The Local course writes its own .env: see labs/local/00-set-up.md.) -# ---------------------------------------------------------------- team card -- +# -------------------------------------------------- from the organizers -- -# Service-account API key (API Key v2) for the hosted Agent Engine API, -# Kafka and Schema Registry. OAuth MCP servers use a separate browser login. +# Your service account's API key (API Key v2). It authenticates the hosted Agent +# Engine API, Kafka and Schema Registry. The MCP server uses a separate browser +# login, in Lab 3. SN_API_KEY= -# Service-account principal, used as the Kafka SASL username. +# Your service account's principal, used as the Kafka SASL username. # Looks like: @.auth.streamnative.cloud SN_SERVICE_ACCOUNT= -# Agent Engine registry endpoint (the External one). Host root only, no /v1. -# Looks like: https:// +# ------------------------------------------------ from your instance -- + +# Your agent workspace's external endpoint, with https:// in front. Host root +# only, no /v1. Looks like: https:// ORCA_BASE_URL= -# Your team's Kafka cluster and its Schema Registry. +# Your Kafka cluster's external endpoint, with its port (host:9093), and its +# schema registry's external endpoint, with https:// in front. KAFKA_BOOTSTRAP_SERVERS= SCHEMA_REGISTRY_URL= -# StreamNative MCP server for your SQL Workspace (the agent's data tools). +# Your SQL workspace's MCP route (the agent's data tools): +# https://mcp.streamnative.cloud/mcp/x//sqlworkspace.compute.streamnative.io/ SN_MCP_URL= # MCP authentication: oauth (default) or static_bearer for API-key MCP servers. @@ -39,15 +45,14 @@ SN_MCP_OAUTH_SCOPE="openid profile email offline_access" # ------------------------------------------------------------- your choices -- -# The Kafka topic name. Injectors and doctor read this value. -# The SQL files in sql/cloud/ use the default below: edit their quoted -# "avro." source name to match this value before running them in -# SQL Workspace. +# The Kafka topic you create and load in Lab 0. The seeder, the injector and the +# doctor read this value. Your SQL catalog imports the topic as a source with the +# same name, and the SQL files in sql/cloud/ use the default below: if you change +# it, change the quoted source name in them too. SQL Workspace does not read .env. LOGIN_TOPIC=security.login_events -# The model your agent runs on (served by the event's AI gateway). +# The model your agent runs on (served by the Agent Engine's AI gateway). ORCA_MODEL=claude-sonnet-4-6 -# Names your agent, environment, and vault so teammates sharing a workspace -# don't collide. Defaults to your OS user name. +# Names your agent, environment, and vault. Defaults to your OS user name. PARTICIPANT= diff --git a/README.md b/README.md index 3a6a605..be52d1f 100644 --- a/README.md +++ b/README.md @@ -25,18 +25,18 @@ The same five labs, on two stacks. | | [Cloud course](labs/cloud/README.md) | [Local course](labs/local/README.md) | |---|---|---| -| Runs on | StreamNative Cloud: your team's Kafka cluster, SQL Workspace, and a hosted Agent Engine | Your laptop: [Ursa for Kafka](https://openlakestream.org/docs/ursa-for-kafka), [RisingWave](https://risingwave.com), and the Orca Agent Engine (`ork local`) | -| You need | A team card, handed out at the hackathon | Docker and an Anthropic API key | -| Time | About 30 minutes | About 45 minutes, plus image downloads | +| Runs on | StreamNative Cloud: your own instance, with a Kafka cluster, a SQL workspace, and an agent workspace | Your laptop: [Ursa for Kafka](https://openlakestream.org/docs/ursa-for-kafka), [RisingWave](https://risingwave.com), and the Orca Agent Engine (`ork local`) | +| You need | A StreamNative Cloud login with your own instance, from the hackathon organizers | Docker and an Anthropic API key | +| Time | About 40 minutes | About 45 minutes, plus image downloads | | Start | [Lab 0: Set up](labs/cloud/00-set-up.md) | [Lab 0: Set up](labs/local/00-set-up.md) | At the hackathon, take the Cloud course: see -[Before you arrive](docs/before-you-arrive.md). Without a team card, or to see -every part run on your own machine, take the Local course. +[Before you arrive](docs/before-you-arrive.md). Without a StreamNative Cloud +instance, or to see every part run on your own machine, take the Local course. ## Pick your path -The agent steps work three ways. Pick one; a teammate can pick another. +The agent steps work three ways. Pick one. - **CLI**: the [`ork`](https://github.com/orca-ae/orca-cli) command line - **Python**: the [`runorca`](https://pypi.org/project/runorca/) SDK diff --git a/cli/env.sh b/cli/env.sh index 846a640..ff639ed 100644 --- a/cli/env.sh +++ b/cli/env.sh @@ -50,7 +50,7 @@ hello_setup_hint() { if [ "$(hello_trim "${TUTORIAL_STACK:-}")" = local ]; then printf '%s' "Run local/write-env.sh in the repo root to write .env again (Local course, Lab 0)." else - printf '%s' "Copy .env.cloud.example to .env in the repo root and fill it in from your team card, or run local/write-env.sh for the Local course." + printf '%s' "Copy .env.cloud.example to .env in the repo root and fill it in from your StreamNative Cloud instance (Cloud course, Lab 0), or run local/write-env.sh for the Local course." fi } @@ -84,7 +84,7 @@ hello_setup() { hello_die "jq is not installed. Install it (brew install jq, apt install jq, or winget install jqlang.jq) and try again." hello_load_dotenv - # `cloud`: your team card on StreamNative Cloud. `local`: the stack on your laptop. + # `cloud`: your instance on StreamNative Cloud. `local`: the stack on your laptop. HELLO_STACK=$(hello_trim "${TUTORIAL_STACK:-cloud}") case "$HELLO_STACK" in cloud | local) ;; diff --git a/cli/tests/run.sh b/cli/tests/run.sh index 5d4f3e3..c2e90e9 100755 --- a/cli/tests/run.sh +++ b/cli/tests/run.sh @@ -133,7 +133,7 @@ test_missing_team_card() { run "" l1_hello.sh check "exits 1 without a team card" [ "$STATUS" -eq 1 ] check "names every missing variable, and both ways to get an .env" \ - err_has "Missing ORCA_BASE_URL, ORCA_MODEL, SN_API_KEY. Copy .env.cloud.example to .env in the repo root and fill it in from your team card, or run local/write-env.sh for the Local course." + err_has "Missing ORCA_BASE_URL, ORCA_MODEL, SN_API_KEY. Copy .env.cloud.example to .env in the repo root and fill it in from your StreamNative Cloud instance (Cloud course, Lab 0), or run local/write-env.sh for the Local course." check "runs no ork command" [ ! -s "$FAKE_ORK_DIR/calls.log" ] } diff --git a/docs/before-you-arrive.md b/docs/before-you-arrive.md index 31f5fae..f3ce0c1 100644 --- a/docs/before-you-arrive.md +++ b/docs/before-you-arrive.md @@ -1,14 +1,15 @@ # Before you arrive -Ten minutes at home saves thirty at the event. This page prepares your laptop -for the [Cloud course](../labs/cloud/README.md), the one you take at the -hackathon. Pick **one** path; your teammate can pick a different one. +Twenty minutes at home saves an hour at the event. This page prepares your +laptop and your StreamNative Cloud instance for the +[Cloud course](../labs/cloud/README.md), the one you take at the hackathon. Pick +**one** path. | Path | Install | |---|---| | **Python** | Python 3.11 or newer | | **TypeScript** | Node.js 20 or newer | -| **CLI** | Python 3.11+ *or* Node.js 20+, for two helper scripts (the doctor and the data injector) | +| **CLI** | Python 3.11+ *or* Node.js 20+, for three helper scripts (the doctor, the seeder, and the data injector) | Everyone also needs: @@ -25,6 +26,9 @@ Everyone also needs: server when a choice is needed. The browser flow stores tokens directly in the vault. - [`jq`](https://jqlang.org/download/), for the checks in every lab. +- [`snctl`](https://docs.streamnative.io/tools/cli/snctl/snctl-overview) (the + StreamNative Cloud CLI): `brew install streamnative/streamnative/snctl`. Lab 0 + uses it to read your instance's addresses and to create your topic. ## 1. Get the code @@ -49,7 +53,7 @@ npm install ``` **CLI**: install `ork` and `jq`, then set up Python or TypeScript as above for -the doctor and the injector. +the doctor, the seeder, and the injector. ## 3. Check your laptop @@ -58,12 +62,39 @@ python doctor.py --offline # Python or CLI path npm run doctor -- --offline # TypeScript path ``` -Every line should say `PASS`. You'll get your **team card** (your credentials -and endpoints) at the event; [Lab 0](../labs/cloud/00-set-up.md) starts there. +Every line should say `PASS`. + +## 4. Set up your instance + +The organizers add you to the hackathon organization on StreamNative Cloud, give +you an **instance** of your own, and make a **service account** in it. They give +you its name and its **API key**: keep the key to yourself. + +In the StreamNative Cloud console, create three things in your instance, in the +region the organizers name: + +- a **Kafka cluster** (Serverless), +- an **agent workspace**, +- a **SQL workspace** that imports your Kafka cluster. + +Then log `snctl` in and check that all three are there: + +```bash +snctl config init +snctl auth login # opens your browser +snctl config set --organization # the hackathon organization's id, o-... +snctl get kafkaclusters -o custom-columns=NAME:.metadata.name,INSTANCE:.spec.instanceName +snctl get workspaces -o custom-columns=NAME:.metadata.name,INSTANCE:.spec.instanceName +snctl get sqlcatalogs -o custom-columns=NAME:.metadata.name,KAFKA_CLUSTER:.spec.sourceRef.name,SQL_WORKSPACE:.spec.workspaceRef.name +``` + +The first two lists have a row for your instance, and the SQL catalog list has a +row that names your Kafka cluster and your SQL workspace. +[Lab 0](../labs/cloud/00-set-up.md) reads their addresses into `.env`. ## Want to try it tonight? The [Local course](../labs/local/README.md) is the same five labs on your own -laptop, with no team card: Ursa for Kafka, RisingWave, and the Orca Agent -Engine in Docker. It needs Docker and an Anthropic API key, and downloads about -5 GB of images, so start it on a good connection. +laptop, with nothing on StreamNative Cloud: Ursa for Kafka, RisingWave, and the +Orca Agent Engine in Docker. It needs Docker and an Anthropic API key, and +downloads about 5 GB of images, so start it on a good connection. diff --git a/docs/tutor.md b/docs/tutor.md index b61e3c5..58b6693 100644 --- a/docs/tutor.md +++ b/docs/tutor.md @@ -49,7 +49,7 @@ last row works in any agent that can read a file. Started with no request, the tutor asks three things: ```text -1. Course: Cloud (a team card, StreamNative Cloud) or Local (everything on your laptop)? +1. Course: Cloud (your own instance on StreamNative Cloud) or Local (everything on your laptop)? 2. Path: CLI, Python, or TypeScript? 3. What now: start at Lab 0, resume at a lab, quiz me on a lab, or check my setup? ``` diff --git a/labs/README.md b/labs/README.md index ef0670a..e6e4cbd 100644 --- a/labs/README.md +++ b/labs/README.md @@ -4,13 +4,13 @@ Two courses teach the same five labs on two stacks. Pick one. | | [Cloud course](cloud/README.md) | [Local course](local/README.md) | |---|---|---| -| Runs on | StreamNative Cloud: your team's Kafka cluster, SQL Workspace, and a hosted Agent Engine | Your laptop: Ursa for Kafka, RisingWave, and the Orca Agent Engine | -| You need | A team card, handed out at the hackathon | Docker and an Anthropic API key | -| Time | About 30 minutes | About 45 minutes, plus the image downloads | -| Take it when | You are at the event | You have no team card, or you want to see every part run | +| Runs on | StreamNative Cloud: your own instance, with a Kafka cluster, a SQL workspace, and an agent workspace | Your laptop: Ursa for Kafka, RisingWave, and the Orca Agent Engine | +| You need | A StreamNative Cloud login with your own instance, from the hackathon organizers | Docker and an Anthropic API key | +| Time | About 40 minutes | About 45 minutes, plus the image downloads | +| Take it when | You are at the event | You have no StreamNative Cloud instance, or you want to see every part run | Both courses use the same three paths for the agent steps. Pick one path and stay -on it; a teammate can pick another. +on it. - **CLI**: the [`ork`](https://github.com/orca-ae/orca-cli) command line - **Python**: the [`runorca`](https://pypi.org/project/runorca/) SDK diff --git a/labs/cloud/00-set-up.md b/labs/cloud/00-set-up.md index ccd193d..4350353 100644 --- a/labs/cloud/00-set-up.md +++ b/labs/cloud/00-set-up.md @@ -1,13 +1,30 @@ # Lab 0: Set up -**Cloud course** · 3 minutes, plus 5 on your own · CLI, Python, or TypeScript +**Cloud course** · 10 minutes, plus 5 on your own · CLI, Python, or TypeScript -You put your team card in `.env` and run the doctor. When this lab is done, you -know that the Agent Engine, Kafka, and Schema Registry on your card all answer. +You fill in `.env` from your instance on StreamNative Cloud, load the login +stream into your Kafka cluster, and run the doctor. When this lab is done, the +Agent Engine, Kafka, and Schema Registry in your instance all answer, and your +topic holds 246 logins. ## Before you start -- You have your **team card** from the organizers. +- You can log in to StreamNative Cloud. The organizers added you to the + hackathon organization, made you an **instance** of your own, and gave you a + **service account** in it: its name and its **API key**. +- In your instance there is a **Kafka cluster**, an **agent workspace**, and a + **SQL workspace** that imports your Kafka cluster. If you have not created + them yet, see [Before you arrive](../../docs/before-you-arrive.md). +- You have [`snctl`](https://docs.streamnative.io/tools/cli/snctl/snctl-overview), + logged in to the hackathon organization: + + ```bash + brew install streamnative/streamnative/snctl + snctl config init + snctl auth login # opens your browser + snctl config set --organization # the hackathon organization's id, o-... + ``` + - You cloned this repository and opened a terminal in it. The terminal runs `bash`: on Windows that is WSL or Git Bash, on every path, because the checks are `bash` commands. @@ -17,7 +34,7 @@ know that the Agent Engine, Kafka, and Schema Registry on your card all answer. ## Step 1: Install your path -Pick **one** path. Your teammate can pick a different one. +Pick **one** path. **Python** (3.11 or newer) @@ -36,7 +53,7 @@ npm install ``` **CLI**: `ork` and `jq` are all the labs need. Set up Python or TypeScript as -above too: the doctor and the data injector come from one of them. +above too: the doctor, the seeder, and the data injector come from one of them. ### Check @@ -57,26 +74,57 @@ PASS jq found All good: you're ready. ``` -## Step 2: Paste your team card +## Step 2: Fill in `.env` from your instance -Open a second terminal at the repository root. Copy the template, then paste the -values from your team card into `.env`: +Open a second terminal at the repository root and copy the template: ```bash cp .env.cloud.example .env ``` -`.env` is git-ignored. It holds your team's key: do not commit it or paste it -anywhere. +`.env` is git-ignored. It will hold your key: do not commit it or paste it +anywhere. Put your service account in it first, as the organizers gave it to +you: + +```text +SN_API_KEY= +SN_SERVICE_ACCOUNT=@.auth.streamnative.cloud +``` + +Then find the names of your three resources. Each list shows the instance a +resource belongs to; the SQL catalog ties your Kafka cluster to your SQL +workspace: + +```bash +snctl get kafkaclusters -o custom-columns=NAME:.metadata.name,DISPLAY:.spec.displayName,INSTANCE:.spec.instanceName +snctl get workspaces -o custom-columns=NAME:.metadata.name,DISPLAY:.spec.displayName,INSTANCE:.spec.instanceName +snctl get sqlcatalogs -o custom-columns=NAME:.metadata.name,KAFKA_CLUSTER:.spec.sourceRef.name,SQL_WORKSPACE:.spec.workspaceRef.name +``` + +Each prints something like this. Use the rows with your instance: -`SN_API_KEY` authenticates the hosted Agent Engine, Kafka, and Schema Registry. -The MCP server uses a separate browser login, in Lab 3: keep `SN_MCP_AUTH=oauth` -and leave `SN_MCP_OAUTH_ISSUER` empty. +```text +NAME DISPLAY INSTANCE +c-abc1234 ana-kafka ana +``` + +Now ask for each address, and write it into `.env`: + +| `.env` line | Command | Write it as | +|---|---|---| +| `ORCA_BASE_URL` | `snctl get workspace -o jsonpath='{.status.serviceEndpoints[?(@.type=="external")].dnsName}'` | `https://` and the host | +| `KAFKA_BOOTSTRAP_SERVERS` | `snctl get kafkacluster -o jsonpath='{.status.serviceEndpoints[?(@.type=="external")].dnsName}'` | as printed, with `:9093` | +| `SCHEMA_REGISTRY_URL` | `snctl get schemaregistry -o jsonpath='{.status.serviceEndpoints[?(@.type=="external")].dnsName}'` | `https://` and the host | +| `SN_MCP_URL` | (no command: it is built from two names) | `https://mcp.streamnative.cloud/mcp/x//sqlworkspace.compute.streamnative.io/` | + +The schema registry has the same name as its Kafka cluster. Leave the other +lines as they are: the MCP server uses a separate browser login in Lab 3, so +`SN_MCP_AUTH=oauth` stays, and `SN_MCP_OAUTH_ISSUER` stays empty. ### Check -One authenticated read of your Agent Engine. It prints `true` when the endpoint -and the key on your card are accepted. +One authenticated read of your Agent Engine. It prints `true` when the address +and the key in `.env` are accepted. ```bash ./lab-ork agent list -o json | jq -e 'has("data")' @@ -85,16 +133,50 @@ and the key on your card are accepted. Before you filled in `.env`, the same command says `Missing ORCA_BASE_URL, SN_API_KEY` instead: the two values it needs. -## Step 3: Run the doctor +## Step 3: Load the login stream + +Your Kafka cluster is new and empty. Create the login topic in it. The first +`snctl kafka` command opens your browser for one more login: + +```bash +snctl context use --instance --kafka-cluster +snctl kafka admin topics create security.login_events --partitions 1 +``` + +```text +Topic 'security.login_events' created successfully with 1 partitions and replication factor 1 +``` + +Then load the stream into it. In your path's folder: + +| Python | TypeScript | CLI | +|---|---|---| +| `python seed.py` | `npm run seed` | `(cd ../python && .venv/bin/python seed.py)` or `(cd ../typescript && npm run seed)` | + +```text +Loaded 246 logins for 91 accounts into security.login_events. +``` + +The seeder replays [`data/login_events.jsonl`](../../data/login_events.jsonl): +synthetic logins at a fictional bank, with their timestamps moved to now. It +also registers the topic's Avro schema, which SQL Workspace needs in Lab 2. +Run it again, from any path, and it refuses to load a second copy, which would +double every count in Lab 2: + +```text +security.login_events already holds 246 events, so it is seeded. To load another copy anyway: python seed.py --force +``` + +### Check -In your path's folder: +Run the doctor. It checks your laptop, then each service in `.env`, and a failed +check prints its fix on the next line. | Python | TypeScript | CLI | |---|---|---| | `python doctor.py` | `npm run doctor` | `(cd ../python && .venv/bin/python doctor.py)` or `(cd ../typescript && npm run doctor)` | -The doctor checks your laptop, then each service on your card. A failed check -prints its fix on the next line. After the lines from step 1, it prints: +After the lines from step 1, it prints: ```text PASS .env cloud stack, participant: ana @@ -110,16 +192,8 @@ You're ready. 1 check(s) wait for a later lab. ``` `WAIT` is not a failure. The MCP server needs a login that only Lab 3 can do. - -### Check - -The last line of the doctor says you are ready, and no line says `FAIL`. - -```bash -(cd python && .venv/bin/python doctor.py) | tail -n 1 -``` - -On the TypeScript path, use `npm --prefix typescript run doctor | tail -n 1`. +Before this step, the `Kafka` and `Schema Registry` lines fail: the topic and +its schema are not there yet. Still failing after two tries? Raise your hand, or see [Troubleshooting](troubleshooting.md). @@ -130,7 +204,7 @@ Still failing after two tries? Raise your hand, or see - A. Fix it now: the doctor has to print only `PASS`. - B. Nothing yet: Lab 3 does the browser login this check waits for. -- C. Ask for a new team card. +- C. Ask the organizers for a new API key.
Answer @@ -141,7 +215,21 @@ A real problem prints `FAIL` and its fix.
-**2. What does `./lab-ork` add to `ork`?** +**2. Where does `ORCA_BASE_URL` come from?** + +- A. Your Kafka cluster's address. +- B. Your agent workspace's external endpoint. +- C. The MCP server. + +
+Answer + +**B.** The Agent Engine runs in your agent workspace. The Kafka cluster gives +you `KAFKA_BOOTSTRAP_SERVERS`, and the SQL workspace gives you `SN_MCP_URL`. + +
+ +**3. What does `./lab-ork` add to `ork`?** - A. It is a different CLI with its own commands. - B. It points `ork` at your Agent Engine with the key from `.env`, and fills in the ids your scripts saved. @@ -169,6 +257,8 @@ The doctor ends on the "ready" line again. (cd python && .venv/bin/python doctor.py) | tail -n 1 ``` +On the TypeScript path, use `npm --prefix typescript run doctor | tail -n 1`. +
Solution @@ -187,8 +277,10 @@ Engine check, tells you the exact value to use, and ends with ## Recap -- `.env` in the repository root is your team card. It is git-ignored. -- The doctor checks each service on the card and prints the fix for a failure. +- `.env` holds the addresses of your three resources and your service account's + key. `snctl` reads the addresses from your instance. `.env` is git-ignored. +- Your Kafka cluster starts empty: you created the topic and loaded it. +- The doctor checks each service and prints the fix for a failure. - `./lab-ork` is how you look at your Agent Engine from the terminal. ## What's next diff --git a/labs/cloud/02-streaming-sql.md b/labs/cloud/02-streaming-sql.md index a1d198b..c9027e9 100644 --- a/labs/cloud/02-streaming-sql.md +++ b/labs/cloud/02-streaming-sql.md @@ -1,6 +1,6 @@ # Lab 2: Hello, streaming SQL -**Cloud course** · 8 minutes, plus 5 on your own · SQL Workspace, in the StreamNative Cloud console +**Cloud course** · 8 minutes, plus 5 on your own · SQL Workspace, in the StreamNative Cloud console or with `psql` You turn the login topic into a materialized view that keeps a running summary per account. When this lab is done, there is a view your agent can read in @@ -9,21 +9,29 @@ Lab 3, and a table it can write to in Lab 4. ## Before you start - You finished [Lab 1](01-hello-agent.md). -- In the StreamNative Cloud console, open **SQL Workspace**, select the hackathon - workspace, and pick your team's database. Use a new query tab for each step. -- **Align the SQL with your `.env` first.** The default Kafka topic is - `security.login_events`, and SQL Workspace exposes its Avro source as - `"avro.security.login_events"`. The injector and the doctor read `LOGIN_TOPIC` - from `.env`, but SQL Workspace does not: the SQL files and the examples below - contain a fixed source name. Check `LOGIN_TOPIC`, then replace - `"avro.security.login_events"` with `"avro."` in - [`sql/cloud/01_explore.sql`](../../sql/cloud/01_explore.sql) and +- Open your SQL workspace. In the StreamNative Cloud console, open **SQL + Workspace**, select your SQL workspace, and pick the database named after your + SQL catalog (Lab 0, step 2). Use a new query tab for each step. +- If the console cannot open the database yet, use `psql` from the repository + root instead. Look up your SQL workspace's address, then connect as `root` + with your API key as the password: + + ```bash + snctl get sqlworkspace -o jsonpath='{.status.endpoints[?(@.type=="sqlgateway/pgwire")].url}' + export PGPASSWORD="$(sed -n 's/^SN_API_KEY=//p' .env)" + psql "postgresql://root@:4567/?sslmode=require" + ``` + + No `psql` on your laptop? Docker has one: + `docker run --rm -it -e PGPASSWORD postgres:16-alpine psql "postgresql://root@:4567/?sslmode=require"`. +- **The source is named after the topic.** Your SQL catalog imported the topic + `security.login_events` as the source `"security.login_events"`. If + `LOGIN_TOPIC` in `.env` is something else, use that name instead, in + [`sql/cloud/01_explore.sql`](../../sql/cloud/01_explore.sql), [`sql/cloud/02_login_failures.sql`](../../sql/cloud/02_login_failures.sql), and - in any query copied from this page. For example, - `LOGIN_TOPIC=security.team07_logins` requires `FROM "avro.security.team07_logins"`. - Keep the double quotes around the entire source name, and confirm that SQL - Workspace imported that topic as an Avro source. Keep the `login_failures` - view name: Labs 3 and 4 query that view. + the queries on this page: SQL Workspace does not read `.env`. Keep the double + quotes around the whole name, and keep the `login_failures` view name: Labs 3 + and 4 query that view. ## Step 1: Peek at the stream @@ -32,7 +40,7 @@ one login attempt. The topic name contains dots, so it is double-quoted. ```sql SELECT event_time, account_id, ip_address, result, failure_reason -FROM "avro.security.login_events" +FROM "security.login_events" ORDER BY event_time DESC LIMIT 20; ``` @@ -40,14 +48,14 @@ LIMIT 20; ### Check The query returns 20 rows, newest first, and `result` is `SUCCESS` or `FAILURE`. -This counts what the source holds: it returns a number greater than zero. +This counts what the source holds: `246`, the logins you loaded in Lab 0. ```sql -SELECT count(*) AS logins FROM "avro.security.login_events"; +SELECT count(*) AS logins FROM "security.login_events"; ``` -If it says `relation "avro.security.login_events" does not exist`, you are in -the wrong database or the source name does not match your topic: see +If it says `table or source not found: security.login_events`, you are in the +wrong database, or the source name does not match your topic: see [Troubleshooting](troubleshooting.md). ## Step 2: Turn the stream into context @@ -63,10 +71,13 @@ SELECT COUNT(*) FILTER (WHERE result = 'SUCCESS') AS successful_logins, COUNT(DISTINCT ip_address) AS distinct_ips, MAX(event_time) AS last_seen -FROM "avro.security.login_events" +FROM "security.login_events" GROUP BY account_id; ``` +It prints a `NOTICE` about snapshot backfill along with +`CREATE_MATERIALIZED_VIEW`. The notice is expected. + A materialized view is maintained incrementally: every new login updates the counts within seconds. There is no batch job to schedule and nothing to refresh. That makes it good agent context: always current, and cheap to read. @@ -90,7 +101,7 @@ FROM login_failures WHERE account_id = 'acct_0042' AND failed_logins >= 5 AND successful_logins >= 1; ``` -Before the step, the same query fails: `login_failures` does not exist. +Before the step, the same query fails: `table or source not found: login_failures`. ## Step 3: Make room for the agent's decisions @@ -129,7 +140,7 @@ as events arrive, so reading it is a cheap lookup and the result is current.
-**2. Why is the source written as `"avro.security.login_events"`, in double quotes?** +**2. Why is the source written as `"security.login_events"`, in double quotes?** - A. The name contains dots, and without quotes each dot would separate a schema from a name. - B. Double quotes make the query case-insensitive. diff --git a/labs/cloud/03-live-context.md b/labs/cloud/03-live-context.md index c61354b..c8bcb33 100644 --- a/labs/cloud/03-live-context.md +++ b/labs/cloud/03-live-context.md @@ -9,7 +9,7 @@ materialized view, and its answer changes when the stream does. ## Before you start - You finished [Lab 2](02-streaming-sql.md): `login_failures` exists in your - team's database. + SQL workspace's database. - One terminal is in your path's folder, a second one is at the repository root. - `ork` v0.6.0 or newer is installed. All three paths use it for the first MCP login. diff --git a/labs/cloud/04-act-with-approval.md b/labs/cloud/04-act-with-approval.md index 871fc36..b81db9f 100644 --- a/labs/cloud/04-act-with-approval.md +++ b/labs/cloud/04-act-with-approval.md @@ -10,7 +10,10 @@ yes, and has not flagged another because you said no. - You finished [Lab 3](03-live-context.md): the agent reads `login_failures`, and the browser login for the MCP server is done. -- `flagged_accounts` exists in your team's database (Lab 2, step 3). +- `flagged_accounts` exists in your SQL workspace's database (Lab 2, step 3). +- Your SQL workspace's MCP access is read-write; the organizers set this up. + Read-only access offers the agent no tool that writes: see + [Troubleshooting](troubleshooting.md). - One terminal is in your path's folder, a second one is at the repository root. ## Step 1: Ask the agent to act, and approve @@ -73,11 +76,12 @@ Ask again (Enter to quit): Now flag acct_0042 as well. [approve?] The agent wants to run sql_workspace_insert_rows with: ... Allow it? [y/N] n -[agent] The insert was denied by a human reviewer, so acct_0042 has not been flagged. +[error] The human reviewer denied this action. +[agent] The human reviewer **denied** this flag. `acct_0042` has **not** been added to `flagged_accounts`, and I will not retry. ``` -The agent is told a human denied the insert, and it does not retry. Press Enter -to end the conversation. +The `[error]` line is the result the agent got back for its tool call: a human +denied it. The agent does not retry. Press Enter to end the conversation. ### Check diff --git a/labs/cloud/README.md b/labs/cloud/README.md index 5a6cd92..bbe4b05 100644 --- a/labs/cloud/README.md +++ b/labs/cloud/README.md @@ -1,10 +1,10 @@ # Cloud course -**Data + Agent Hackathon: hello world, on StreamNative Cloud · about 30 minutes** +**Data + Agent Hackathon: hello world, on StreamNative Cloud · about 40 minutes** You build an agent whose context is a live Kafka stream, kept fresh by streaming -SQL, and that asks a human before it acts. Everything runs on the environment on -your team card. +SQL, and that asks a human before it acts. Everything runs in your own instance +on StreamNative Cloud. ## The story @@ -24,7 +24,7 @@ flowchart LR | Lab | Time | Where | You | The idea | |---|---|---|---|---| -| [0. Set up](00-set-up.md) | 3 min | terminal | Fill in `.env`, run the doctor | Check service access before you build on it | +| [0. Set up](00-set-up.md) | 10 min | terminal | Fill in `.env` from your instance, load the topic, run the doctor | Check service access before you build on it | | [1. Hello, agent](01-hello-agent.md) | 5 min | CLI / Python / TS | Create an agent and chat | Agent, environment, session, events | | [2. Hello, streaming SQL](02-streaming-sql.md) | 8 min | SQL Workspace | Build a materialized view over the topic | Context that keeps itself fresh | | [3. Agent + live context](03-live-context.md) | 9 min | CLI / Python / TS | Give the agent SQL tools, inject new data | The answer changes with the data | @@ -35,9 +35,12 @@ your own. ## What you need -- Your **team card** from the organizers. Your team's Kafka cluster already - holds the login stream. -- One path installed, plus `ork` and `jq`: see +- A login to StreamNative Cloud, in the hackathon organization, with an + **instance** of your own and a **service account** in it (its name and API + key). The organizers set these up. +- In your instance, a **Kafka cluster**, an **agent workspace**, and a **SQL + workspace** that imports the Kafka cluster. +- One path installed, plus `ork`, `jq`, and `snctl`. All of this is in [Before you arrive](../../docs/before-you-arrive.md). Start with [Lab 0: Set up](00-set-up.md). If something goes wrong, see @@ -45,19 +48,24 @@ Start with [Lab 0: Set up](00-set-up.md). If something goes wrong, see [The labs](../README.md), and a coding agent can [tutor you through the course](../../docs/tutor.md). -No team card? Take the [Local course](../local/README.md): the same labs, on -your laptop. +No StreamNative Cloud instance? Take the [Local course](../local/README.md): +the same labs, on your laptop. ## What was run -These pages were rewritten on 2 October 2026, and the checks were added then. - -- **Run, on the Agent Engine of the [Local course](../local/README.md)** (`ork` - 0.6.0, which serves the same API): the check commands of Labs 1, 3 and 4 that - read your agent, its versions, its tools and their policies, and its sessions. -- **Not run for this revision: anything that needs a team card.** That is the - doctor against StreamNative Cloud, SQL Workspace (Lab 2), the browser login - and the vault (Lab 3), the hosted SQL tools (Labs 3 and 4), and every step - where the model answers. For those steps the labs show what the scripts are - written to print and what the checks are written to select, not a recording of - a run. +These pages were rewritten on 2 October 2026 and run against one test instance +on StreamNative Cloud that night, with `snctl` 1.8.0 and `ork` 0.6.0. The Kafka +cluster was Serverless; the SQL workspace ran RisingWave 3.1.0-alpha. + +- **Lab 0**: every `snctl` lookup, the topic, the seeder, and the doctor, on the + Python path. On the TypeScript path, the doctor, and the seeder against the + topic once it was loaded. +- **Lab 2**: every statement and check, through `psql`. The console was not + used. +- **Labs 1, 3 and 4**: every step and check, on all three paths, with the model + answering: the injected account showing up, one insert approved and one + denied. On the CLI path the first Lab 4 run gave up after five minutes: the + Agent Engine acted on the approval eight minutes after it was given. The + second run passed. +- **Not run on this course**: the "Try it yourself" tasks, which were run on + the Local course, and the clean-up at the end of Lab 4. diff --git a/labs/cloud/troubleshooting.md b/labs/cloud/troubleshooting.md index 8fac950..f326055 100644 --- a/labs/cloud/troubleshooting.md +++ b/labs/cloud/troubleshooting.md @@ -1,7 +1,7 @@ # Troubleshooting: Cloud course -Run the doctor first. It checks each service on your team card and prints the -fix for what fails. +Run the doctor first. It checks each service in your `.env` and prints the fix +for what fails. | Python or CLI | TypeScript | |---|---| @@ -15,16 +15,24 @@ Still stuck after two tries? Raise your hand. |---|---| | Doctor: `Agent Engine HTTP 401/403` | The key was rejected. A key created before its permissions must be re-created: ask a facilitator. | | Doctor: `Kafka ... authentication` | `SN_SERVICE_ACCOUNT` must be the full principal, `@.auth.streamnative.cloud`; `SN_API_KEY` is the raw key. | +| Doctor: `Kafka security.login_events: not found` | The topic is not there yet. Create it and load it: Lab 0, step 3. | +| Doctor: `Schema Registry ... not found` | The schema is registered when you load the topic: Lab 0, step 3. | | Doctor: `WAIT MCP OAuth` | Not a failure. Lab 3 does the browser login; run the doctor again after it. | -| The login topic isn't listed in SQL Workspace | Only topics with a registered Avro schema appear. Ask a facilitator. | -| `relation "avro.security.login_events" does not exist` | Select your team's database, and update the quoted Avro source in both Lab 2 SQL files to match `LOGIN_TOPIC` in `.env`. | -| The agent can't find `login_failures` | Create the view in your team's database (Lab 2, step 2); the agent looks it up there. | -| `[error]` lines from MCP tools in Lab 3 | Check `SN_MCP_URL` and `SN_MCP_AUTH`, finish the OAuth login, then run the doctor again. | +| `snctl kafka ...`: `organization, instance and pulsar cluster are required` | Point `snctl` at your cluster first: `snctl context use --instance --kafka-cluster ` (Lab 0, step 3). | +| `snctl get ...` lists nothing, or another organization's resources | Set the hackathon organization: `snctl config set --organization `. | +| The SQL workspace never gets ready (`snctl get sqlworkspace ` shows `does not enable SQLWorkspace`) | It was created in a region that has no SQL workspaces. Ask a facilitator which region to use. | +| The console cannot open your SQL workspace's database | Use `psql` instead: Lab 2, "Before you start". | +| `table or source not found: security.login_events` | Pick the database named after your SQL catalog. If your `LOGIN_TOPIC` is not `security.login_events`, use your topic's name in both Lab 2 SQL files: the source is named exactly after the topic. | +| The agent can't find `login_failures` | Create the view in your SQL workspace's database (Lab 2, step 2); the agent looks it up there. | +| `[error] connection_unavailable: sql connection unavailable` from every SQL tool | The MCP server reached your SQL workspace but could not log in to its RisingWave as you. If `psql` works (Lab 2), the SQL workspace runs a RisingWave version that does not accept that login: ask a facilitator to update it. | +| Lab 4: the agent says `sql_workspace_insert_rows` is not available or the session is read-only, and `Allow it? [y/N]` never appears | Your SQL workspace's MCP access is read-only. Lab 4 needs it read-write, which a facilitator sets for you (in the console: your SQL workspace, Settings, MCP). Then run the Lab 4 script again. | +| Other `[error]` lines from MCP tools in Lab 3 | Check `SN_MCP_URL` (your SQL workspace's route, Lab 0, step 2) and `SN_MCP_AUTH`, finish the OAuth login, then run the doctor again. | | OAuth issuer mismatch / unsupported client authentication | Use `ork` v0.6.0 or newer and leave `SN_MCP_OAUTH_ISSUER` empty for StreamNative discovery. An explicit issuer must match an advertised authorization server. `--oauth-allow-issuer-mismatch` is only for trusted servers whose metadata issuer crosses registrable domains; StreamNative does not need it. | -| `Cannot reach the Agent Engine` | `ORCA_BASE_URL` must be the host root from your card, with no `/v1`. | +| `Cannot reach the Agent Engine` | `ORCA_BASE_URL` must be `https://` and your agent workspace's external endpoint, with no `/v1`: Lab 0, step 2. | | The agent answers from memory instead of querying | Ask again, "check the view first". The system prompt tells it to always query. | | `./lab-ork` says `No session_id yet` | The lab step that creates it has not run on this stack. Run the lab's script first. | | A script seems stuck at an approval | The session is waiting for you. Answer the `Allow it? [y/N]` prompt, or press Ctrl-C and run the lab script again: it starts a fresh session. | +| Nothing happens for minutes after you answer `y` or `n` (on the CLI path: `The agent did not finish its turn within 300s.`) | The Agent Engine has your decision but has not acted on it yet. Run the Lab 4 script again: it starts a fresh session. A `y` you already gave can still be applied later, so the account you approved may be in `flagged_accounts` already. | ## The MCP login in Lab 3 @@ -53,5 +61,5 @@ For the StreamNative SQL Workspace MCP server, keep `SN_MCP_AUTH=oauth`, leave - Agent Engine: run the cleanup script of your path (`./cleanup.sh`, `python cleanup.py`, or `npm run cleanup`). - SQL: run [`sql/cloud/99_reset.sql`](../../sql/cloud/99_reset.sql) in your - team's database. + SQL workspace's database. - Kafka: injected `acct_9…` events stay in the topic. They are harmless. diff --git a/local/tests/run.sh b/local/tests/run.sh index 2b9f316..54e9019 100755 --- a/local/tests/run.sh +++ b/local/tests/run.sh @@ -151,7 +151,8 @@ test_write_env_never_overwrites_a_team_card() { printf 'SN_API_KEY=team-key\nORCA_BASE_URL=https://ws.example.com\n' >"$R/.env" run write-env.sh check "exits 1 on a cloud .env" [ "$STATUS" -eq 1 ] - check "leaves the team card as it was" env_is SN_API_KEY team-key + check "leaves the cloud .env as it was" env_is SN_API_KEY team-key + check "says which course that .env is for" err_has ".env is set up for the Cloud course" check "says how to keep both" err_has "mv .env .env.cloud" } diff --git a/local/write-env.sh b/local/write-env.sh index 867678c..661cd78 100755 --- a/local/write-env.sh +++ b/local/write-env.sh @@ -6,7 +6,7 @@ # local/write-env.sh # # Run it again whenever local/engine.sh has started a fresh engine. It keeps your -# PARTICIPANT and ORCA_MODEL, and it never overwrites a team card. +# PARTICIPANT and ORCA_MODEL, and it never overwrites a Cloud course .env. set -euo pipefail # shellcheck source=lib.sh . "$(dirname "${BASH_SOURCE[0]}")/lib.sh" @@ -21,7 +21,7 @@ participant="" model=claude-sonnet-4-6 if [ -f "$ENV_FILE" ]; then [ "$(env_value TUTORIAL_STACK "$ENV_FILE")" = local ] || - die ".env holds a team card (the Cloud course), and this would replace it. + die ".env is set up for the Cloud course, and this would replace it. To keep it, move it aside first: mv .env .env.cloud (both names are git-ignored)" participant=$(env_value PARTICIPANT "$ENV_FILE") kept=$(env_value ORCA_MODEL "$ENV_FILE") diff --git a/python/common.py b/python/common.py index 161452a..0141799 100644 --- a/python/common.py +++ b/python/common.py @@ -50,7 +50,7 @@ def __getitem__(self, name: str) -> str: @property def stack(self) -> str: - """`cloud`: your team card on StreamNative Cloud. `local`: the stack on your laptop.""" + """`cloud`: your instance on StreamNative Cloud. `local`: the stack on your laptop.""" stack = self.values.get("TUTORIAL_STACK", "cloud") if stack not in STACKS: raise ConfigError("TUTORIAL_STACK must be cloud or local.") @@ -67,8 +67,8 @@ def setup_hint(values: Mapping[str, str]) -> str: if values.get("TUTORIAL_STACK") == "local": return "Run local/write-env.sh in the repo root to write .env again (Local course, Lab 0)." return ( - "Copy .env.cloud.example to .env in the repo root and fill it in from your team card, " - "or run local/write-env.sh for the Local course." + "Copy .env.cloud.example to .env in the repo root and fill it in from your StreamNative Cloud instance " + "(Cloud course, Lab 0), or run local/write-env.sh for the Local course." ) @@ -188,12 +188,12 @@ def schema_registry_config(config: Config) -> dict[str, str]: def orca_client(config: Config) -> Orca: - """Use a Registry workspace key locally, or the team card's hosted Bearer key.""" + """Use a Registry workspace key locally, or the service account's API key as a Bearer token on StreamNative Cloud.""" if key := config.values.get("ORCA_API_KEY"): return Orca(base_url=config["ORCA_BASE_URL"], api_key=None, default_headers={"x-api-key": key}, timeout=600) if key := config.values.get("SN_API_KEY"): return Orca(base_url=config["ORCA_BASE_URL"], api_key=key, timeout=600) - raise ConfigError("Set ORCA_API_KEY for ork local, or SN_API_KEY from your team card.") + raise ConfigError("Set ORCA_API_KEY for ork local, or SN_API_KEY (your service account's API key) for StreamNative Cloud.") def ensure_environment(client: Any, state: State, name: str) -> str: diff --git a/python/doctor.py b/python/doctor.py index a30cbaa..55e2479 100644 --- a/python/doctor.py +++ b/python/doctor.py @@ -4,7 +4,7 @@ python doctor.py --offline # laptop only (run this before the event) python doctor.py --agent-only # laptop + Agent Engine (enough for Lab 1) -It checks the stack your .env is for: your team card on StreamNative Cloud, or +It checks the stack your .env is for: your instance on StreamNative Cloud, or the stack on your laptop. Every failed check prints the fix. """ @@ -59,7 +59,7 @@ def check_orca_base_url(url: str) -> Check: parts = urlsplit(url.strip()) local_http = parts.scheme == "http" and parts.hostname in ("localhost", "127.0.0.1", "::1") if not parts.netloc or (parts.scheme != "https" and not local_http): - return Check("ORCA_BASE_URL", False, url, "Use the https:// registry endpoint from your team card, or http://127.0.0.1:8080 for ork local.") + return Check("ORCA_BASE_URL", False, url, "Use your agent workspace's https:// external endpoint (Cloud course, Lab 0), or http://127.0.0.1:8080 for ork local.") root = f"{parts.scheme}://{parts.netloc}" if parts.path.rstrip("/"): return Check("ORCA_BASE_URL", False, url, f"Use the host root only: ORCA_BASE_URL={root}") @@ -132,12 +132,16 @@ def kafka_hint(error: str, stack: str = "cloud") -> str: if "authorization" in text: return "Your key logs in but may not use this topic: its rolebinding is missing. Ask a facilitator." if any(word in text for word in ("resolve", "transport", "timed out", "connect")): - return "Cannot reach Kafka. Check KAFKA_BOOTSTRAP_SERVERS (host:port from your team card) and your network." + return "Cannot reach Kafka. Check KAFKA_BOOTSTRAP_SERVERS (your Kafka cluster's host:port, Cloud course, Lab 0) and your network." + if "not found" in text: + return "The login topic is not there yet. Create it and load it: Cloud course, Lab 0." return "See the error above, or ask a facilitator." def schema_registry_hint(error: str, stack: str = "cloud") -> str: if stack != "local": + if "not found" in error.lower(): + return "The schema is registered when you load the topic: python seed.py (Cloud course, Lab 0)." return "Check SCHEMA_REGISTRY_URL; your key may lack Schema Registry read access." if "not found" in error.lower(): return "The schema is registered when you seed the topic: python seed.py (Local course, Lab 0)." @@ -198,9 +202,9 @@ def probe_orca(config: Any) -> Check: if err.status_code in (401, 403) and local: fix = "The key in .env does not match the running stack. Run local/write-env.sh; if it still fails, start over with local/down.sh --reset." elif err.status_code in (401, 403): - fix = "The Agent Engine rejected the key. For ork local use its generated workspace key as ORCA_API_KEY; for a team card check SN_API_KEY and its rolebinding." + fix = "The Agent Engine rejected the key. For ork local use its generated workspace key as ORCA_API_KEY; on StreamNative Cloud check SN_API_KEY and its rolebinding." elif err.status_code == 404: - fix = "ORCA_BASE_URL is not an Agent Engine registry: copy the registry endpoint from your team card." + fix = "ORCA_BASE_URL is not an Agent Engine registry: use your agent workspace's external endpoint (Cloud course, Lab 0)." else: fix = "Ask a facilitator." return Check("Agent Engine", False, f"HTTP {err.status_code}", fix) diff --git a/python/seed.py b/python/seed.py index 3f6c3b8..42e39d0 100644 --- a/python/seed.py +++ b/python/seed.py @@ -1,4 +1,4 @@ -"""Load the login stream into the topic on your laptop (Local course, Lab 0). +"""Load the login stream into your topic (Lab 0, in either course). Replays data/login_events.jsonl: 246 synthetic logins at a fictional bank, with their timestamps moved to now. One of the accounts in it is under attack. @@ -38,7 +38,7 @@ def events_in_topic(watermarks: Iterable[tuple[int, int]]) -> int: return sum(high - low for low, high in watermarks) -def count_existing(consumer: Any, topic: str) -> int: +def count_existing(consumer: Any, topic: str, stack: str = "local") -> int: """How many events the topic holds already. Stops the script when it cannot tell.""" from confluent_kafka import KafkaException, TopicPartition @@ -46,7 +46,8 @@ def count_existing(consumer: Any, topic: str) -> int: # Listing every topic avoids a metadata request that could create a missing one. metadata = consumer.list_topics(timeout=15).topics.get(topic) if metadata is None or metadata.error is not None: - raise SystemExit(f"The topic {topic} does not exist yet. Create it first: Local course, Lab 0.") + course = "Local course" if stack == "local" else "Cloud course" + raise SystemExit(f"The topic {topic} does not exist yet. Create it first: {course}, Lab 0.") return events_in_topic(consumer.get_watermark_offsets(TopicPartition(topic, p), timeout=15) for p in metadata.partitions) except KafkaException as err: reason = err.args[0].str() if err.args else str(err) @@ -59,11 +60,9 @@ def main() -> None: from confluent_kafka import Consumer config = load_config(["KAFKA_BOOTSTRAP_SERVERS", "SCHEMA_REGISTRY_URL", "LOGIN_TOPIC"]) - if config.stack != "local": - raise SystemExit("seed.py loads the topic on your laptop (Local course). Your team's cluster already holds the login stream.") topic = config["LOGIN_TOPIC"] - existing = count_existing(Consumer({**kafka_client_config(config), "group.id": "hello-seed"}), topic) + existing = count_existing(Consumer({**kafka_client_config(config), "group.id": "hello-seed"}), topic, config.stack) if existing and "--force" not in sys.argv[1:]: raise SystemExit(f"{topic} already holds {existing} events, so it is seeded. To load another copy anyway: python seed.py --force") diff --git a/python/tests/test_doctor.py b/python/tests/test_doctor.py index 4f27595..f191c0c 100644 --- a/python/tests/test_doctor.py +++ b/python/tests/test_doctor.py @@ -2,6 +2,7 @@ import re from pathlib import Path +from types import SimpleNamespace import pytest @@ -15,6 +16,7 @@ kafka_hint, mcp_headers, parse_mcp_response, + probe_orca, required_for, schema_registry_hint, summarize, @@ -217,13 +219,58 @@ def test_an_unreachable_local_schema_registry_points_at_the_streaming_stack(): @pytest.mark.parametrize("error", ["[Errno 61] Connection refused", "Unauthorized (HTTP status code 401, SR code 401)"]) -def test_cloud_schema_registry_errors_point_at_the_team_card(error): +def test_cloud_schema_registry_errors_point_at_the_schema_registry_url(error): hint = schema_registry_hint(error) assert "SCHEMA_REGISTRY_URL" in hint assert "compose" not in hint +# On StreamNative Cloud each participant creates and loads their own topic. + + +def test_a_cloud_topic_that_is_not_there_yet_points_at_lab_0(): + hint = kafka_hint("security.login_events: not found", "cloud") + + assert "Cloud course, Lab 0" in hint + assert "facilitator" not in hint + + +def test_a_cloud_schema_that_is_not_registered_yet_points_at_the_seeder(): + # What StreamNative Cloud's registry says: it names the subject with its namespace. + hint = schema_registry_hint("Subject 'public/default/security.login_events-value' not found. (HTTP status code 404, SR code 40401)") + + assert "python seed.py" in hint + assert "Cloud course, Lab 0" in hint + + +@pytest.mark.parametrize("hint", [ + kafka_hint("Failed to resolve kafka.example.com:9093", "cloud"), + check_orca_base_url("http://ws.example.com").fix, +]) +def test_cloud_fixes_name_your_instance_not_a_team_card(hint): + assert "team card" not in hint + + +@pytest.mark.parametrize("status", [401, 404]) +def test_the_agent_engine_fixes_for_the_cloud_course_name_no_team_card(monkeypatch, status): + import httpx2 + from orca import APIStatusError + + request = httpx2.Request("GET", "https://ws.example.com/v1/agents") + error = APIStatusError("refused", response=httpx2.Response(status, request=request), body=None) + + class Agents: + def list(self, limit): + raise error + + monkeypatch.setattr("common.orca_client", lambda config: SimpleNamespace(agents=Agents())) + check = probe_orca(SimpleNamespace(stack="cloud")) + + assert not check.ok + assert "team card" not in check.fix + + def test_the_mcp_probe_sends_a_bearer_token_only_when_it_has_one(): assert mcp_headers("the-token")["Authorization"] == "Bearer the-token" assert "Authorization" not in mcp_headers(None) diff --git a/python/tests/test_seed.py b/python/tests/test_seed.py index c9323be..f0e81e9 100644 --- a/python/tests/test_seed.py +++ b/python/tests/test_seed.py @@ -1,4 +1,4 @@ -"""seed.py: the login stream the Local course loads into your topic.""" +"""seed.py: the login stream each course loads into your topic in Lab 0.""" import ipaddress import json @@ -11,6 +11,8 @@ from fastavro import parse_schema from fastavro.validation import validate +import seed +from common import Config from seed import count_existing, events_in_topic, load_events, rebase_events TOPIC = "security.login_events" @@ -132,6 +134,40 @@ def test_a_missing_topic_stops_the_seed_and_points_at_lab_0(): assert consumer.closed +@pytest.mark.parametrize("stack, course", [("local", "Local course"), ("cloud", "Cloud course")]) +def test_a_missing_topic_names_the_course_whose_lab_0_creates_it(stack, course): + with pytest.raises(SystemExit) as stop: + count_existing(FakeConsumer(), TOPIC, stack) + + assert f"{course}, Lab 0" in str(stop.value) + + +def test_the_seed_loads_your_own_cluster_on_the_cloud_stack_too(monkeypatch, capsys): + # On StreamNative Cloud every participant has their own instance and cluster, + # and nobody has loaded it for them. + cloud = Config( + values={ + "SN_API_KEY": "the-api-key", + "SN_SERVICE_ACCOUNT": "test@o-test.auth.streamnative.cloud", + "KAFKA_BOOTSTRAP_SERVERS": "kafka.example.com:9093", + "SCHEMA_REGISTRY_URL": "https://sr.example.com", + "LOGIN_TOPIC": TOPIC, + }, + participant="jane", + ) + written: list = [] + monkeypatch.setattr(seed, "load_config", lambda names: cloud) + monkeypatch.setattr(seed, "count_existing", lambda consumer, topic, stack="local": 0) + monkeypatch.setattr(seed, "login_producer", lambda config, schema: "producer") + monkeypatch.setattr(seed, "publish", lambda producer, topic, events: written.extend(events) or []) + monkeypatch.setattr("confluent_kafka.Consumer", lambda settings: "consumer") + + seed.main() + + assert len(written) == 246 + assert f"Loaded 246 logins for 91 accounts into {TOPIC}." in capsys.readouterr().out + + def test_an_unreachable_broker_stops_the_seed_with_the_reason_and_a_next_step(): consumer = FakeConsumer(unreachable=True) diff --git a/python/tests/test_stack.py b/python/tests/test_stack.py index c8182b9..6101ef1 100644 --- a/python/tests/test_stack.py +++ b/python/tests/test_stack.py @@ -112,6 +112,16 @@ def test_without_a_stack_the_hint_names_both_ways_to_get_an_env_file(): assert "local/write-env.sh" in str(err.value) +def test_the_cloud_hint_sends_you_to_your_own_instance_not_to_a_team_card(): + # Every participant reads the values from their own StreamNative Cloud + # instance (Cloud course, Lab 0); nobody hands out a card any more. + with pytest.raises(ConfigError) as err: + load_config(["ORCA_BASE_URL"], env={}) + + assert "Cloud course, Lab 0" in str(err.value) + assert "team card" not in str(err.value) + + def test_on_the_local_stack_the_hint_is_to_write_the_env_file_again(): with pytest.raises(ConfigError) as err: load_config(["ORCA_BASE_URL"], env={"TUTORIAL_STACK": "local"}) diff --git a/skills/data-agent-tutor/SKILL.md b/skills/data-agent-tutor/SKILL.md index e91c2cb..d1b097c 100644 --- a/skills/data-agent-tutor/SKILL.md +++ b/skills/data-agent-tutor/SKILL.md @@ -43,7 +43,7 @@ pages, in `labs/cloud/` and in `labs/local/`: `03-live-context.md` · `04-act-with-approval.md` · `troubleshooting.md` - **They named no course, or no lab and no task**: reply with only this menu. - 1. Course: **Cloud** (a team card, StreamNative Cloud) or **Local** (everything on your laptop)? + 1. Course: **Cloud** (your own instance on StreamNative Cloud) or **Local** (everything on your laptop)? 2. Path: **CLI**, **Python**, or **TypeScript**? 3. What now: **start** at Lab 0, **resume** at a lab, **quiz me** on a lab, or **check my setup**? - **They pasted an error or a `FAIL` line**: give the fix the troubleshooting diff --git a/sql/cloud/01_explore.sql b/sql/cloud/01_explore.sql index 0b8d1a2..a04f815 100644 --- a/sql/cloud/01_explore.sql +++ b/sql/cloud/01_explore.sql @@ -1,11 +1,12 @@ --- L2 step 1: peek at the live login stream in your team's Kafka cluster. --- Before running: match the quoted "avro.security.login_events" source below --- to LOGIN_TOPIC in .env ("avro."). SQL Workspace does not load .env. +-- Lab 2, step 1: peek at the live login stream in your Kafka cluster. +-- Before running: the source below is named exactly after the topic. If your +-- LOGIN_TOPIC in .env is not security.login_events, use your topic's name here. +-- SQL Workspace does not read .env. -- --- SQL Workspace imported the topic as a source table. Its name contains dots, --- so always wrap it in double quotes. +-- The SQL catalog imported the topic as a source. Its name contains dots, so +-- always wrap it in double quotes. SELECT event_time, account_id, ip_address, result, failure_reason -FROM "avro.security.login_events" +FROM "security.login_events" ORDER BY event_time DESC LIMIT 20; diff --git a/sql/cloud/02_login_failures.sql b/sql/cloud/02_login_failures.sql index 1106ed5..2607c42 100644 --- a/sql/cloud/02_login_failures.sql +++ b/sql/cloud/02_login_failures.sql @@ -1,6 +1,7 @@ --- L2 step 2: turn the stream into always-fresh context for the agent. --- Before running: match the quoted "avro.security.login_events" source below --- to LOGIN_TOPIC in .env ("avro."). SQL Workspace does not load .env. +-- Lab 2, step 2: turn the stream into always-fresh context for the agent. +-- Before running: the source below is named exactly after the topic. If your +-- LOGIN_TOPIC in .env is not security.login_events, use your topic's name here. +-- SQL Workspace does not read .env. -- -- A materialized view is maintained incrementally: every new login event -- updates the counts within seconds. No batch job, no refresh. @@ -12,7 +13,7 @@ SELECT COUNT(*) FILTER (WHERE result = 'SUCCESS') AS successful_logins, COUNT(DISTINCT ip_address) AS distinct_ips, MAX(event_time) AS last_seen -FROM "avro.security.login_events" +FROM "security.login_events" GROUP BY account_id; -- Check it: the accounts with the most failed logins. diff --git a/sql/cloud/03_flagged_accounts.sql b/sql/cloud/03_flagged_accounts.sql index d2ae9a4..856360c 100644 --- a/sql/cloud/03_flagged_accounts.sql +++ b/sql/cloud/03_flagged_accounts.sql @@ -1,4 +1,4 @@ --- L2 step 3: the table your agent will write to in L4 (with your approval). +-- Lab 2, step 3: the table your agent will write to in Lab 4 (with your approval). CREATE TABLE flagged_accounts ( account_id VARCHAR PRIMARY KEY, diff --git a/sql/cloud/99_reset.sql b/sql/cloud/99_reset.sql index 2ef5b54..2317754 100644 --- a/sql/cloud/99_reset.sql +++ b/sql/cloud/99_reset.sql @@ -1,4 +1,4 @@ --- Start L2 over: drop what the tutorial created. The Kafka topic is untouched. +-- Start Lab 2 over: drop what the lab created. The Kafka topic and its source are untouched. DROP TABLE IF EXISTS flagged_accounts; DROP MATERIALIZED VIEW IF EXISTS login_failures; diff --git a/typescript/src/common.ts b/typescript/src/common.ts index f10dcb0..43836e9 100644 --- a/typescript/src/common.ts +++ b/typescript/src/common.ts @@ -62,7 +62,7 @@ export class Config { return this.#values[name]; } - /** `cloud`: your team card on StreamNative Cloud. `local`: the stack on your laptop. */ + /** `cloud`: your instance on StreamNative Cloud. `local`: the stack on your laptop. */ get stack(): Stack { const stack = this.#values.TUTORIAL_STACK ?? 'cloud'; if (stack !== 'cloud' && stack !== 'local') throw new ConfigError('TUTORIAL_STACK must be cloud or local.'); @@ -81,8 +81,8 @@ export function setupHint(config: Config): string { return 'Run local/write-env.sh in the repo root to write .env again (Local course, Lab 0).'; } return ( - 'Copy .env.cloud.example to .env in the repo root and fill it in from your team card, ' + - 'or run local/write-env.sh for the Local course.' + 'Copy .env.cloud.example to .env in the repo root and fill it in from your StreamNative Cloud instance ' + + '(Cloud course, Lab 0), or run local/write-env.sh for the Local course.' ); } @@ -292,14 +292,14 @@ export interface Client { }; } -/** Use a Registry workspace key locally, or the team card's hosted Bearer key. */ +/** Use a Registry workspace key locally, or the service account's API key as a Bearer token on StreamNative Cloud. */ export function orcaClient(config: Config): Orca { const baseURL = config.get('ORCA_BASE_URL'); if (config.has('ORCA_API_KEY')) { return new Orca({ baseURL, apiKey: null, defaultHeaders: { 'x-api-key': config.get('ORCA_API_KEY') }, timeout: 600_000 }); } if (config.has('SN_API_KEY')) return new Orca({ baseURL, apiKey: config.get('SN_API_KEY'), timeout: 600_000 }); - throw new ConfigError('Set ORCA_API_KEY for ork local, or SN_API_KEY from your team card.'); + throw new ConfigError("Set ORCA_API_KEY for ork local, or SN_API_KEY (your service account's API key) for StreamNative Cloud."); } /** The sandbox your sessions run in. Created once, then reused. */ diff --git a/typescript/src/doctor.ts b/typescript/src/doctor.ts index 5248afe..0b5d049 100644 --- a/typescript/src/doctor.ts +++ b/typescript/src/doctor.ts @@ -5,7 +5,7 @@ * npm run doctor -- --offline # laptop only (run this before the event) * npm run doctor -- --agent-only # laptop + Agent Engine (enough for Lab 1) * - * It checks the stack your .env is for: your team card on StreamNative Cloud, or + * It checks the stack your .env is for: your instance on StreamNative Cloud, or * the stack on your laptop. Every failed check prints the fix. */ @@ -62,7 +62,7 @@ export function checkOrcaBaseUrl(url: string): Check { } const localHttp = parsed?.protocol === 'http:' && ['localhost', '127.0.0.1', '[::1]'].includes(parsed.hostname); if (!parsed || !parsed.host || (parsed.protocol !== 'https:' && !localHttp)) { - return check('ORCA_BASE_URL', false, url, 'Use the https:// registry endpoint from your team card, or http://127.0.0.1:8080 for ork local.'); + return check('ORCA_BASE_URL', false, url, "Use your agent workspace's https:// external endpoint (Cloud course, Lab 0), or http://127.0.0.1:8080 for ork local."); } const root = `${parsed.protocol}//${parsed.host}`; if (parsed.pathname.replace(/\/+$/, '')) { @@ -147,14 +147,19 @@ export function kafkaHint(error: string, stack: Stack = 'cloud'): string { return 'Your key logs in but may not use this topic: its rolebinding is missing. Ask a facilitator.'; } if (unreachable) { - return 'Cannot reach Kafka. Check KAFKA_BOOTSTRAP_SERVERS (host:port from your team card) and your network.'; + return "Cannot reach Kafka. Check KAFKA_BOOTSTRAP_SERVERS (your Kafka cluster's host:port, Cloud course, Lab 0) and your network."; } + if (text.includes('not found')) return 'The login topic is not there yet. Create it and load it: Cloud course, Lab 0.'; return 'See the error above, or ask a facilitator.'; } export function schemaRegistryHint(error: string, stack: Stack = 'cloud'): string { - if (stack !== 'local') return 'Check SCHEMA_REGISTRY_URL; your key may lack Schema Registry read access.'; - if (error.toLowerCase().includes('not found')) return 'The schema is registered when you seed the topic: npm run seed (Local course, Lab 0).'; + const notFound = error.toLowerCase().includes('not found'); + if (stack !== 'local') { + if (notFound) return 'The schema is registered when you load the topic: npm run seed (Cloud course, Lab 0).'; + return 'Check SCHEMA_REGISTRY_URL; your key may lack Schema Registry read access.'; + } + if (notFound) return 'The schema is registered when you seed the topic: npm run seed (Local course, Lab 0).'; return `Cannot reach Schema Registry on your laptop. ${START_LOCAL_STACK}`; } @@ -202,6 +207,20 @@ function checkDocker(): Check { return check('docker', found, found ? 'found' : 'not found', found ? '' : 'The Local course runs in Docker: install Docker Desktop, or Docker Engine with Compose v2.'); } +/** What to do when the Agent Engine answers with an HTTP error status. */ +export function agentEngineFix(status: number, stack: Stack): string { + if ((status === 401 || status === 403) && stack === 'local') { + return 'The key in .env does not match the running stack. Run local/write-env.sh; if it still fails, start over with local/down.sh --reset.'; + } + if (status === 401 || status === 403) { + return 'The Agent Engine rejected the key. For ork local use its generated workspace key as ORCA_API_KEY; on StreamNative Cloud check SN_API_KEY and its rolebinding.'; + } + if (status === 404) { + return "ORCA_BASE_URL is not an Agent Engine registry: use your agent workspace's external endpoint (Cloud course, Lab 0)."; + } + return 'Ask a facilitator.'; +} + async function probeOrca(config: Config): Promise { const { APIConnectionError, APIError } = await import('@runorca/orca-sdk'); const { orcaClient } = await import('./common.js'); @@ -215,15 +234,7 @@ async function probeOrca(config: Config): Promise { return check('Agent Engine', false, err.message, fix); } if (err instanceof APIError && err.status) { - let fix = 'Ask a facilitator.'; - if ((err.status === 401 || err.status === 403) && local) { - fix = 'The key in .env does not match the running stack. Run local/write-env.sh; if it still fails, start over with local/down.sh --reset.'; - } else if (err.status === 401 || err.status === 403) { - fix = 'The Agent Engine rejected the key. For ork local use its generated workspace key as ORCA_API_KEY; for a team card check SN_API_KEY and its rolebinding.'; - } else if (err.status === 404) { - fix = 'ORCA_BASE_URL is not an Agent Engine registry: copy the registry endpoint from your team card.'; - } - return check('Agent Engine', false, `HTTP ${err.status}`, fix); + return check('Agent Engine', false, `HTTP ${err.status}`, agentEngineFix(err.status, config.stack)); } throw err; } diff --git a/typescript/src/seed.ts b/typescript/src/seed.ts index 2e68871..fcc362c 100644 --- a/typescript/src/seed.ts +++ b/typescript/src/seed.ts @@ -1,5 +1,5 @@ /** - * Load the login stream into the topic on your laptop (Local course, Lab 0). + * Load the login stream into your topic (Lab 0, in either course). * * Replays data/login_events.jsonl: 246 synthetic logins at a fictional bank, with * their timestamps moved to now. One of the accounts in it is under attack. @@ -14,7 +14,7 @@ import { pathToFileURL } from 'node:url'; import { Kafka, logLevel } from 'kafkajs'; -import { REPO_ROOT, kafkaClientConfig, loadConfig, runMain, type Config } from './common.js'; +import { REPO_ROOT, kafkaClientConfig, loadConfig, runMain, type Config, type Stack } from './common.js'; import { loginProducer, publish, type LoginEvent } from './inject.js'; const EVENTS_FILE = join(REPO_ROOT, 'data', 'login_events.jsonl'); @@ -51,12 +51,13 @@ export interface TopicAdmin { } /** How many events the topic holds already. Stops the script when it cannot tell. */ -export async function countExisting(admin: TopicAdmin, topic: string): Promise { +export async function countExisting(admin: TopicAdmin, topic: string, stack: Stack = 'local'): Promise { try { await admin.connect(); // Listing every topic avoids a metadata request that could create a missing one. if (!(await admin.listTopics()).includes(topic)) { - throw new SeedError(`The topic ${topic} does not exist yet. Create it first: Local course, Lab 0.`); + const course = stack === 'local' ? 'Local course' : 'Cloud course'; + throw new SeedError(`The topic ${topic} does not exist yet. Create it first: ${course}, Lab 0.`); } return eventsInTopic(await admin.fetchTopicOffsets(topic)); } catch (err) { @@ -72,24 +73,32 @@ function topicAdmin(config: Config): TopicAdmin { return new Kafka({ clientId: 'hello-seed', ...kafkaClientConfig(config), logLevel: logLevel.NOTHING, retry: { retries: 2 } }).admin(); } -async function main(): Promise { - const config = loadConfig(['KAFKA_BOOTSTRAP_SERVERS', 'SCHEMA_REGISTRY_URL', 'LOGIN_TOPIC']); - if (config.stack !== 'local') { - throw new SeedError("`npm run seed` loads the topic on your laptop (Local course). Your team's cluster already holds the login stream."); - } +/** Load the events into the topic, unless it holds some already. Returns what to tell the participant. */ +export async function seed( + config: Config, + admin: TopicAdmin, + write: (topic: string, events: LoginEvent[]) => Promise, + force = false, +): Promise { const topic = config.get('LOGIN_TOPIC'); - - const existing = await countExisting(topicAdmin(config), topic); - if (existing > 0 && !process.argv.slice(2).includes('--force')) { + const existing = await countExisting(admin, topic, config.stack); + if (existing > 0 && !force) { throw new SeedError(`${topic} already holds ${existing} events, so it is seeded. To load another copy anyway: npm run seed -- --force`); } const events = rebaseEvents(loadEvents(), new Date()); - // Registering the schema is what lets RisingWave decode the topic. - const errors = await publish(loginProducer(config, { schema: readFileSync(SCHEMA_FILE, 'utf8') }), topic, events); + const errors = await write(topic, events); if (errors.length > 0) throw new SeedError(`Could not write to ${topic}: ${errors[0]}\nRun \`npm run doctor\` to check your setup.`); const accounts = new Set(events.map((event) => event.account_id)).size; - console.log(`Loaded ${events.length} logins for ${accounts} accounts into ${topic}.`); + return `Loaded ${events.length} logins for ${accounts} accounts into ${topic}.`; +} + +async function main(): Promise { + const config = loadConfig(['KAFKA_BOOTSTRAP_SERVERS', 'SCHEMA_REGISTRY_URL', 'LOGIN_TOPIC']); + // Registering the schema is what lets the streaming database decode the topic. + const write = (topic: string, events: LoginEvent[]) => + publish(loginProducer(config, { schema: readFileSync(SCHEMA_FILE, 'utf8') }), topic, events); + console.log(await seed(config, topicAdmin(config), write, process.argv.slice(2).includes('--force'))); } if (process.argv[1] && import.meta.url === pathToFileURL(process.argv[1]).href) { diff --git a/typescript/test/doctor.test.ts b/typescript/test/doctor.test.ts index 44d46a0..b4bad7c 100644 --- a/typescript/test/doctor.test.ts +++ b/typescript/test/doctor.test.ts @@ -5,6 +5,7 @@ import { existsSync } from 'node:fs'; import { describe, expect, it } from 'vitest'; import { + agentEngineFix, check as newCheck, checkLoginSchema, checkMcpQuery, @@ -235,13 +236,39 @@ describe('the two stacks', () => { expect(hint).not.toContain('seed'); }); - it.each(['fetch failed: connect ECONNREFUSED 10.0.0.1:443', 'HTTP 401: unauthorized'])('cloud schema registry error %s points at the team card', (error) => { + it.each(['fetch failed: connect ECONNREFUSED 10.0.0.1:443', 'HTTP 401: unauthorized'])('cloud schema registry error %s points at SCHEMA_REGISTRY_URL', (error) => { const hint = schemaRegistryHint(error); expect(hint).toContain('SCHEMA_REGISTRY_URL'); expect(hint).not.toContain('compose'); }); + // On StreamNative Cloud each participant creates and loads their own topic. + + it('a cloud topic that is not there yet points at Lab 0', () => { + const hint = kafkaHint('security.login_events: not found', 'cloud'); + + expect(hint).toContain('Cloud course, Lab 0'); + expect(hint).not.toContain('facilitator'); + }); + + it('a cloud schema that is not registered yet points at the seeder', () => { + // What StreamNative Cloud's registry says: it names the subject with its namespace. + const hint = schemaRegistryHint('HTTP 404: {"error_code":40401,"message":"Subject \'public/default/security.login_events-value\' not found."}'); + + expect(hint).toContain('npm run seed'); + expect(hint).toContain('Cloud course, Lab 0'); + }); + + it.each([ + ['an unreachable cluster', kafkaHint('Failed to resolve kafka.example.com:9093', 'cloud')], + ['an http URL', checkOrcaBaseUrl('http://ws.example.com').fix ?? ''], + ['a rejected key', agentEngineFix(401, 'cloud')], + ['a URL that is not a registry', agentEngineFix(404, 'cloud')], + ])('the cloud fix for %s names your instance, not a team card', (_case, hint) => { + expect(hint).not.toContain('team card'); + }); + it('the MCP probe sends a bearer token only when it has one', () => { expect(mcpHeaders('the-token').Authorization).toBe('Bearer the-token'); expect(Object.keys(mcpHeaders())).not.toContain('Authorization'); diff --git a/typescript/test/seed.test.ts b/typescript/test/seed.test.ts index 42a35c5..7135153 100644 --- a/typescript/test/seed.test.ts +++ b/typescript/test/seed.test.ts @@ -1,4 +1,4 @@ -/** seed.ts: the login stream the Local course loads into your topic. */ +/** seed.ts: the login stream each course loads into your topic in Lab 0. */ import { readFileSync } from 'node:fs'; import { BlockList } from 'node:net'; @@ -6,8 +6,9 @@ import { BlockList } from 'node:net'; import avro from 'avsc'; import { describe, expect, it } from 'vitest'; +import { Config } from '../src/common.js'; import type { LoginEvent } from '../src/inject.js'; -import { SeedError, countExisting, eventsInTopic, loadEvents, rebaseEvents } from '../src/seed.js'; +import { SeedError, countExisting, eventsInTopic, loadEvents, rebaseEvents, seed } from '../src/seed.js'; const TOPIC = 'security.login_events'; const SCHEMA = avro.Type.forSchema(JSON.parse(readFileSync(new URL('../../schemas/login_events.avsc', import.meta.url), 'utf8'))); @@ -140,6 +141,50 @@ describe('countExisting', () => { expect((failure as Error).message).toContain('npm run doctor'); expect(admin.disconnected).toBe(true); }); + + it.each([ + ['local', 'Local course'], + ['cloud', 'Cloud course'], + ] as const)('a missing topic names the course whose Lab 0 creates it (%s)', async (stack, course) => { + const failure = await countExisting(fakeAdmin({ offsets: null }), TOPIC, stack).catch((err: unknown) => err); + + expect((failure as Error).message).toContain(`${course}, Lab 0`); + }); +}); + +describe('seed', () => { + it('loads your own cluster on the cloud stack too', async () => { + // On StreamNative Cloud every participant has their own instance and cluster, + // and nobody has loaded it for them. + const cloud = new Config( + { + SN_API_KEY: 'the-api-key', + SN_SERVICE_ACCOUNT: 'test@o-test.auth.streamnative.cloud', + KAFKA_BOOTSTRAP_SERVERS: 'kafka.example.com:9093', + SCHEMA_REGISTRY_URL: 'https://sr.example.com', + LOGIN_TOPIC: TOPIC, + }, + 'jane', + ); + const written: unknown[] = []; + + const done = await seed(cloud, fakeAdmin({ offsets: [{ low: '0', high: '0' }] }), async (_topic, events) => { + written.push(...events); + return []; + }); + + expect(written).toHaveLength(246); + expect(done).toBe(`Loaded 246 logins for 91 accounts into ${TOPIC}.`); + }); + + it('refuses a topic that already holds events, unless forced', async () => { + const failure = await seed(new Config({ TUTORIAL_STACK: 'local', LOGIN_TOPIC: TOPIC }, 'jane'), fakeAdmin(), async () => []).catch( + (err: unknown) => err, + ); + + expect((failure as Error).message).toContain('already holds 246 events'); + expect(await seed(new Config({ TUTORIAL_STACK: 'local', LOGIN_TOPIC: TOPIC }, 'jane'), fakeAdmin(), async () => [], true)).toContain('Loaded 246'); + }); }); describe('eventsInTopic', () => { diff --git a/typescript/test/stack.test.ts b/typescript/test/stack.test.ts index b22d045..d5af3f1 100644 --- a/typescript/test/stack.test.ts +++ b/typescript/test/stack.test.ts @@ -125,6 +125,15 @@ describe('hints', () => { expect(message).toContain('local/write-env.sh'); }); + it('the cloud hint sends you to your own instance, not to a team card', () => { + // Every participant reads the values from their own StreamNative Cloud + // instance (Cloud course, Lab 0); nobody hands out a card any more. + const message = messageOf(() => loadConfig(['ORCA_BASE_URL'], {})); + + expect(message).toContain('Cloud course, Lab 0'); + expect(message).not.toContain('team card'); + }); + it('on the local stack, the hint is to write the .env file again', () => { const message = messageOf(() => loadConfig(['ORCA_BASE_URL'], { TUTORIAL_STACK: 'local' }));