| 1 | //! Codewhale account and BYOK credential commands. |
| 2 | //! |
| 3 | //! This module is deliberately separate from the provider-facing `login` and |
| 4 | //! `auth` commands in `lib.rs`: those configure the local runtime, while this |
| 5 | //! surface signs a CLI profile into the managed Codewhale account and stores |
| 6 | //! provider keys in that account's remote vault. |
| 7 | |
| 8 | use std::io::{self, IsTerminal, Read, Write}; |
| 9 | use std::net::IpAddr; |
| 10 | use std::thread; |
| 11 | use std::time::Duration; |
| 12 | |
| 13 | use anyhow::{Context, Result, anyhow, bail}; |
| 14 | use clap::{Args, Subcommand, ValueEnum}; |
| 15 | use codewhale_config::device_code::DevicePollOutcome; |
| 16 | use codewhale_config::{ConfigStore, ProviderKind}; |
| 17 | use codewhale_secrets::Secrets; |
| 18 | use codewhale_secrets::account::{ |
| 19 | ACCOUNT_API_BASE_ENV as CLOUD_API_BASE_ENV, AccountAuthBundle as AuthBundle, |
| 20 | AccountSessionSnapshot, AccountSessionStore, AccountUser as CloudUser, |
| 21 | DEFAULT_ACCOUNT_API_BASE as DEFAULT_API_BASE, StoredAccountAuth as StoredCloudAuth, |
| 22 | normalize_account_profile as normalized_profile, secure_account_session_secrets, |
| 23 | validate_account_auth_bundle as validate_auth_bundle, |
| 24 | }; |
| 25 | use reqwest::Url; |
| 26 | use serde::{Deserialize, Serialize, de::DeserializeOwned}; |
| 27 | |
| 28 | pub(crate) mod machine; |
| 29 | mod work; |
| 30 | |
| 31 | const MAX_RESPONSE_BYTES: u64 = 256 * 1024; |
| 32 | const MIN_API_KEY_BYTES: usize = 8; |
| 33 | const MAX_API_KEY_BYTES: u64 = 4096; |
| 34 | const MAX_API_KEY_STDIN_BYTES: u64 = MAX_API_KEY_BYTES + 1024; |
| 35 | const MAX_KEY_LABEL_CHARS: usize = 80; |
| 36 | pub(crate) const DEFAULT_LOGIN_TIMEOUT_SECONDS: u64 = 600; |
| 37 | pub(crate) const MAX_LOGIN_TIMEOUT_SECONDS: u64 = 3600; |
| 38 | |
| 39 | #[derive(Debug, Args)] |
| 40 | pub(crate) struct CloudArgs { |
| 41 | /// Codewhale account API origin. HTTPS is required except for loopback HTTP. |
| 42 | #[arg(long, global = true, value_name = "URL")] |
| 43 | api_base: Option<String>, |
| 44 | #[command(subcommand)] |
| 45 | command: CloudCommand, |
| 46 | } |
| 47 | |
| 48 | #[derive(Debug, Subcommand)] |
| 49 | enum CloudCommand { |
| 50 | /// Sign this CLI profile in through the browser device flow. |
| 51 | Login(CloudLoginArgs), |
| 52 | /// Show the signed-in account for this CLI profile. |
| 53 | Status, |
| 54 | /// Remove this profile's local account session and revoke it when reachable. |
| 55 | Logout, |
| 56 | /// Manage provider API keys stored in the signed-in Codewhale account. |
| 57 | /// |
| 58 | /// These are credentials Codewhale presents *to* a model provider. For the |
| 59 | /// machine tokens a customer presents *to* Codewhale, see `api-keys`. |
| 60 | Keys(CloudKeysArgs), |
| 61 | /// Manage Computers in this signed-in Codewhale account. |
| 62 | Computers(CloudComputersArgs), |
| 63 | /// List Projects available to this signed-in Codewhale account. |
| 64 | Projects(CloudProjectsArgs), |
| 65 | /// Inspect GitHub repositories authorized for this account. |
| 66 | Github(CloudGithubArgs), |
| 67 | /// Manage named Agents and their account conversations. |
| 68 | Agents(CloudAgentsArgs), |
| 69 | /// Manage Codewhale account API keys: machine tokens for CI. |
| 70 | #[command(name = "api-keys")] |
| 71 | ApiKeys(machine::ApiKeysArgs), |
| 72 | /// Show the account this CLI authenticates as, preferring a machine key. |
| 73 | Whoami, |
| 74 | /// Check the account's agent-model precondition for machine work. |
| 75 | Agent, |
| 76 | /// Inspect the account document; local settings import is not available yet. |
| 77 | Pull(CloudPullArgs), |
| 78 | /// Push local settings to the account document (never automatic, --dry-run required). |
| 79 | Push(CloudPushArgs), |
| 80 | } |
| 81 | |
| 82 | #[derive(Debug, Args)] |
| 83 | struct CloudLoginArgs { |
| 84 | /// Print the verification URL without trying to open a browser. |
| 85 | #[arg(long, default_value_t = false)] |
| 86 | no_open: bool, |
| 87 | /// Maximum time to wait for browser authorization. |
| 88 | #[arg( |
| 89 | long = "timeout-seconds", |
| 90 | default_value_t = DEFAULT_LOGIN_TIMEOUT_SECONDS, |
| 91 | value_parser = clap::value_parser!(u64).range(1..=MAX_LOGIN_TIMEOUT_SECONDS) |
| 92 | )] |
| 93 | timeout_seconds: u64, |
| 94 | } |
| 95 | |
| 96 | #[derive(Debug, Args)] |
| 97 | struct CloudPullArgs { |
| 98 | /// Inspect the account document without writing local files. |
| 99 | #[arg(long, default_value_t = false)] |
| 100 | dry_run: bool, |
| 101 | } |
| 102 | |
| 103 | #[derive(Debug, Args)] |
| 104 | struct CloudPushArgs { |
| 105 | /// Show what would be pushed without writing the remote document. |
| 106 | #[arg(long, default_value_t = false)] |
| 107 | dry_run: bool, |
| 108 | } |
| 109 | |
| 110 | #[derive(Debug, Args)] |
| 111 | struct CloudKeysArgs { |
| 112 | #[command(subcommand)] |
| 113 | command: CloudKeysCommand, |
| 114 | } |
| 115 | |
| 116 | #[derive(Debug, Args)] |
| 117 | struct CloudComputersArgs { |
| 118 | #[command(subcommand)] |
| 119 | command: CloudComputersCommand, |
| 120 | } |
| 121 | |
| 122 | #[derive(Debug, Args)] |
| 123 | struct CloudAgentsArgs { |
| 124 | #[command(subcommand)] |
| 125 | command: CloudAgentsCommand, |
| 126 | } |
| 127 | |
| 128 | #[derive(Debug, Args)] |
| 129 | struct CloudProjectsArgs { |
| 130 | #[command(subcommand)] |
| 131 | command: CloudProjectsCommand, |
| 132 | } |
| 133 | |
| 134 | #[derive(Debug, Args)] |
| 135 | struct CloudGithubArgs { |
| 136 | #[command(subcommand)] |
| 137 | command: CloudGithubCommand, |
| 138 | } |
| 139 | |
| 140 | #[derive(Debug, Subcommand)] |
| 141 | enum CloudGithubCommand { |
| 142 | /// List repository bindings saved by GitHub App installation. |
| 143 | Bindings { |
| 144 | #[arg(long)] |
| 145 | json: bool, |
| 146 | }, |
| 147 | /// Connect one repository from an installed GitHub App to this account. |
| 148 | /// |
| 149 | /// The account API checks that you own the installation and can read the |
| 150 | /// repository; nothing is written to GitHub. Omit --installation-id when |
| 151 | /// every repository already connected shares one installation. |
| 152 | Bind { |
| 153 | /// Repository as OWNER/REPO. |
| 154 | repo: String, |
| 155 | /// GitHub App installation ID (see `account github bindings --json`). |
| 156 | #[arg(long)] |
| 157 | installation_id: Option<String>, |
| 158 | }, |
| 159 | } |
| 160 | |
| 161 | #[derive(Debug, Subcommand)] |
| 162 | enum CloudProjectsCommand { |
| 163 | /// List account Projects; use an ID when binding an Agent. |
| 164 | List { |
| 165 | #[arg(long)] |
| 166 | json: bool, |
| 167 | }, |
| 168 | /// Create a Project from a GitHub repository already connected to this account. |
| 169 | Create { |
| 170 | name: String, |
| 171 | #[arg(long)] |
| 172 | repo_binding_id: String, |
| 173 | #[arg(long)] |
| 174 | operation_key: String, |
| 175 | }, |
| 176 | } |
| 177 | |
| 178 | #[derive(Debug, Subcommand)] |
| 179 | enum CloudAgentsCommand { |
| 180 | /// List active Agents in this account. |
| 181 | List { |
| 182 | #[arg(long)] |
| 183 | json: bool, |
| 184 | }, |
| 185 | /// Create a named Agent. Reuse the operation key if a response is lost. |
| 186 | Create { |
| 187 | name: String, |
| 188 | /// Assign an existing account Project to this Agent. |
| 189 | #[arg(long)] |
| 190 | project_id: Option<String>, |
| 191 | #[arg(long)] |
| 192 | operation_key: String, |
| 193 | }, |
| 194 | /// Assign an existing Project to an Agent, checking its saved revision. |
| 195 | BindProject { agent: String, project_id: String }, |
| 196 | /// List recent, active conversations for a named Agent or Agent ID. |
| 197 | Threads { |
| 198 | agent: String, |
| 199 | #[arg(long)] |
| 200 | json: bool, |
| 201 | }, |
| 202 | /// Create an Agent conversation with an explicit saved model route. |
| 203 | NewThread { |
| 204 | agent: String, |
| 205 | #[arg(long, default_value = "Main")] |
| 206 | title: String, |
| 207 | #[arg(long, default_value = "deepseek")] |
| 208 | provider: String, |
| 209 | #[arg(long, default_value = "deepseek-flash")] |
| 210 | model: String, |
| 211 | #[arg(long)] |
| 212 | operation_key: String, |
| 213 | }, |
| 214 | /// Send one message to a selected conversation, using its saved model. |
| 215 | Send { |
| 216 | agent: String, |
| 217 | prompt: String, |
| 218 | /// Select a specific conversation; otherwise use its only active one, or a unique Main. |
| 219 | #[arg(long)] |
| 220 | thread: Option<String>, |
| 221 | /// Explicit payer. Only byok_external is available under the current launch policy. |
| 222 | #[arg(long)] |
| 223 | billing_mode: String, |
| 224 | /// Stable message ID for safe retry after an uncertain response. |
| 225 | #[arg(long)] |
| 226 | operation_key: String, |
| 227 | }, |
| 228 | /// Read one turn's answer and status from this account's conversation events. |
| 229 | Result { |
| 230 | agent: String, |
| 231 | thread: String, |
| 232 | turn: String, |
| 233 | /// Resume from a previously printed event sequence for long conversations. |
| 234 | #[arg(long, default_value_t = 0)] |
| 235 | since_seq: u64, |
| 236 | }, |
| 237 | /// Give a named Agent durable repository Work. This records a request; it does not start compute. |
| 238 | /// |
| 239 | /// The Agent reads the message first: a plain instruction becomes Work, a |
| 240 | /// command such as "stop" is applied to its active Work, a correction |
| 241 | /// edits that Work's objective, and a question changes nothing. The |
| 242 | /// output states which happened. If the reply is lost, re-run the same |
| 243 | /// command with the same --message-id; that never creates a second Work for |
| 244 | /// the same instruction. If the message was a stop or a correction, check |
| 245 | /// `work-status` first. |
| 246 | Work { |
| 247 | agent: String, |
| 248 | objective: String, |
| 249 | /// Stable message ID for safe retry after an uncertain response. |
| 250 | #[arg(long)] |
| 251 | message_id: String, |
| 252 | }, |
| 253 | /// Read the account-owned status of a Work request. |
| 254 | WorkStatus { id: String }, |
| 255 | /// Cancel Work and stop its computer; safe to repeat. |
| 256 | /// |
| 257 | /// If the Work still has queued prompts, choose --queue discard or --queue |
| 258 | /// park. An already-finished Work is reported and exits successfully. When |
| 259 | /// the reply is lost the outcome is unknown: run `work-status` before |
| 260 | /// retrying. |
| 261 | WorkCancel { |
| 262 | id: String, |
| 263 | /// What happens to queued prompts: discard them or park them on the Work. |
| 264 | #[arg(long, value_enum)] |
| 265 | queue: Option<WorkQueueChoice>, |
| 266 | /// Short reason recorded with the cancellation. |
| 267 | #[arg(long)] |
| 268 | reason: Option<String>, |
| 269 | }, |
| 270 | /// Read a Work's outcome: state, attempts, evidence, artifacts, draft PR, model route and usage. |
| 271 | /// |
| 272 | /// --json prints the account API's raw result and attempt records. |
| 273 | WorkResult { |
| 274 | id: String, |
| 275 | #[arg(long)] |
| 276 | json: bool, |
| 277 | }, |
| 278 | /// Quote bounded Boat trial computer time for queued Work. Nothing starts. |
| 279 | /// |
| 280 | /// Prints the funding, EU placement and time disclosure plus a short-lived |
| 281 | /// confirmation for `work-launch`. Reuse one --operation-key for the quote |
| 282 | /// and the launch. |
| 283 | WorkQuote { |
| 284 | id: String, |
| 285 | /// Stable launch ID, reused by `work-launch` so a retry cannot start two computers. |
| 286 | #[arg(long)] |
| 287 | operation_key: String, |
| 288 | }, |
| 289 | /// Start quoted Boat trial Work on EU compute (five minutes at most, $0 Codewhale charge). |
| 290 | /// |
| 291 | /// Requires the confirmation printed by `work-quote` and --confirm-eu-compute, |
| 292 | /// which agrees that repository code and Work files run on Boat's EU |
| 293 | /// compute. Re-running with the same --operation-key and --confirmation is |
| 294 | /// safe: it replays the launch instead of starting another computer. |
| 295 | WorkLaunch { |
| 296 | id: String, |
| 297 | #[arg(long)] |
| 298 | operation_key: String, |
| 299 | /// Confirmation printed by `work-quote`; use - to read bounded piped stdin. |
| 300 | #[arg(long)] |
| 301 | confirmation: String, |
| 302 | /// Agree that repository code and Work files run on EU compute. |
| 303 | #[arg(long)] |
| 304 | confirm_eu_compute: bool, |
| 305 | }, |
| 306 | } |
| 307 | |
| 308 | /// What a cancellation does with prompts still waiting behind the Work. |
| 309 | #[derive(Debug, Clone, Copy, PartialEq, Eq, ValueEnum)] |
| 310 | enum WorkQueueChoice { |
| 311 | Discard, |
| 312 | Park, |
| 313 | } |
| 314 | |
| 315 | #[derive(Debug, Subcommand)] |
| 316 | enum CloudComputersCommand { |
| 317 | /// List this account's Computers; --json includes allowance and entitlement data. |
| 318 | List { |
| 319 | #[arg(long)] |
| 320 | json: bool, |
| 321 | }, |
| 322 | /// Save a Computer identity; compute is allocated only when it starts. |
| 323 | Create { |
| 324 | name: String, |
| 325 | /// Try Boat compute in the EU for this Computer. |
| 326 | #[arg(long, requires = "eu_compute_opt_in")] |
| 327 | boat_trial: bool, |
| 328 | /// Confirm that Boat trial code and files run on EU compute. |
| 329 | #[arg(long, requires = "boat_trial")] |
| 330 | eu_compute_opt_in: bool, |
| 331 | }, |
| 332 | /// Show one Computer; --json includes allowance and meter data. |
| 333 | Show { |
| 334 | id: String, |
| 335 | #[arg(long)] |
| 336 | json: bool, |
| 337 | }, |
| 338 | /// Read Boat trial usage receipts for one Computer. |
| 339 | Usage { |
| 340 | id: String, |
| 341 | /// Print the full metering response. |
| 342 | #[arg(long, required = true)] |
| 343 | json: bool, |
| 344 | }, |
| 345 | /// Start a Computer, subject to the account's plan and capacity. |
| 346 | Start { id: String }, |
| 347 | /// Pause a Computer. |
| 348 | Pause { id: String }, |
| 349 | /// Permanently delete a Computer and its provider allocation. |
| 350 | Delete { id: String }, |
| 351 | } |
| 352 | |
| 353 | #[derive(Debug, Subcommand)] |
| 354 | enum CloudKeysCommand { |
| 355 | /// List configured providers without revealing key values. |
| 356 | List, |
| 357 | /// Save a provider key to the signed-in Codewhale account. |
| 358 | Set(CloudKeySetArgs), |
| 359 | /// Remove a provider key from the signed-in Codewhale account. |
| 360 | Remove { |
| 361 | /// Provider id from the account's catalog (`account keys list`). |
| 362 | provider: String, |
| 363 | }, |
| 364 | } |
| 365 | |
| 366 | #[derive(Debug, Args)] |
| 367 | struct CloudKeySetArgs { |
| 368 | /// Provider id from the account's catalog (`account keys list`). |
| 369 | provider: String, |
| 370 | /// Read the key from stdin. Useful for pipes and secret-manager commands. |
| 371 | #[arg(long = "api-key-stdin", conflicts_with = "from_local")] |
| 372 | api_key_stdin: bool, |
| 373 | /// Upload the locally resolved key (config, secret store, then environment). |
| 374 | #[arg(long, conflicts_with = "api_key_stdin")] |
| 375 | from_local: bool, |
| 376 | /// Non-secret label shown beside the stored credential. |
| 377 | #[arg(long, default_value = "Codewhale CLI")] |
| 378 | label: String, |
| 379 | } |
| 380 | |
| 381 | /// One row of the account control plane's public provider catalog. |
| 382 | /// |
| 383 | /// This is untrusted remote data, not a Codewhale-owned enum: the account |
| 384 | /// service adds providers without a CLI release, so the catalog is read as |
| 385 | /// data and every id is re-validated locally before it reaches a URL path. |
| 386 | /// Only the fields this surface actually uses are modeled; unknown fields are |
| 387 | /// ignored rather than being turned into behavior. |
| 388 | #[derive(Debug, Clone, Deserialize)] |
| 389 | #[serde(rename_all = "camelCase")] |
| 390 | struct CatalogProvider { |
| 391 | id: String, |
| 392 | #[serde(default)] |
| 393 | label: String, |
| 394 | /// The runtime provider id this catalog row maps onto, when one exists. |
| 395 | /// `--from-local` uses it to find the local credential; without it the |
| 396 | /// row's own id is tried. |
| 397 | #[serde(default)] |
| 398 | runtime_provider: Option<String>, |
| 399 | /// Model ids the account API lists for this provider. `new-thread` checks |
| 400 | /// its explicit route against this served list, not a compiled table. |
| 401 | #[serde(default)] |
| 402 | models: Vec<String>, |
| 403 | /// False for runtime-only rows a hosted conversation cannot use. |
| 404 | #[serde(default)] |
| 405 | connection_available: Option<bool>, |
| 406 | } |
| 407 | |
| 408 | #[derive(Debug, Deserialize)] |
| 409 | struct ProviderCatalogResponse { |
| 410 | #[serde(default)] |
| 411 | providers: Vec<CatalogProvider>, |
| 412 | } |
| 413 | |
| 414 | #[derive(Debug, Deserialize)] |
| 415 | #[serde(rename_all = "camelCase")] |
| 416 | struct AccountComputer { |
| 417 | id: String, |
| 418 | owner_id: String, |
| 419 | name: String, |
| 420 | region: String, |
| 421 | status: String, |
| 422 | #[serde(default)] |
| 423 | start_queue_reason: String, |
| 424 | } |
| 425 | |
| 426 | #[derive(Deserialize)] |
| 427 | struct ComputerListResponse { |
| 428 | computers: Vec<AccountComputer>, |
| 429 | } |
| 430 | |
| 431 | #[derive(Deserialize)] |
| 432 | struct ComputerResponse { |
| 433 | computer: AccountComputer, |
| 434 | #[serde(default)] |
| 435 | queued: bool, |
| 436 | } |
| 437 | |
| 438 | #[derive(Deserialize)] |
| 439 | #[serde(rename_all = "camelCase")] |
| 440 | struct ComputerDeleteResponse { |
| 441 | deleted: bool, |
| 442 | computer_id: String, |
| 443 | } |
| 444 | |
| 445 | #[derive(Deserialize)] |
| 446 | struct AccountProject { |
| 447 | id: String, |
| 448 | name: String, |
| 449 | #[serde(rename = "defaultRepoProvider", default)] |
| 450 | default_repo_provider: String, |
| 451 | #[serde(rename = "defaultRepo", default)] |
| 452 | default_repo: String, |
| 453 | } |
| 454 | |
| 455 | #[derive(Deserialize)] |
| 456 | struct ProjectListResponse { |
| 457 | projects: Vec<AccountProject>, |
| 458 | } |
| 459 | |
| 460 | #[derive(Deserialize)] |
| 461 | struct ProjectResponse { |
| 462 | project: AccountProject, |
| 463 | } |
| 464 | |
| 465 | #[derive(Deserialize)] |
| 466 | #[serde(rename_all = "camelCase")] |
| 467 | struct AccountGitHubBinding { |
| 468 | id: String, |
| 469 | provider: String, |
| 470 | repo: String, |
| 471 | status: String, |
| 472 | #[serde(default)] |
| 473 | installation_id: String, |
| 474 | } |
| 475 | |
| 476 | #[derive(Deserialize)] |
| 477 | struct GitHubBindingListResponse { |
| 478 | bindings: Vec<AccountGitHubBinding>, |
| 479 | } |
| 480 | |
| 481 | #[derive(Serialize)] |
| 482 | struct ComputerCreateRequest<'a> { |
| 483 | name: &'a str, |
| 484 | #[serde(skip_serializing_if = "Option::is_none")] |
| 485 | provider: Option<&'static str>, |
| 486 | #[serde(rename = "boatEuComputeOptIn", skip_serializing_if = "Option::is_none")] |
| 487 | boat_eu_compute_opt_in: Option<bool>, |
| 488 | } |
| 489 | |
| 490 | #[derive(Clone, Deserialize)] |
| 491 | #[serde(rename_all = "camelCase")] |
| 492 | struct AccountAgent { |
| 493 | id: String, |
| 494 | name: String, |
| 495 | status: String, |
| 496 | #[serde(default)] |
| 497 | project_id: String, |
| 498 | revision: u64, |
| 499 | } |
| 500 | |
| 501 | #[derive(Deserialize)] |
| 502 | struct AgentListResponse { |
| 503 | agents: Vec<AccountAgent>, |
| 504 | } |
| 505 | |
| 506 | #[derive(Deserialize)] |
| 507 | struct AgentResponse { |
| 508 | agent: AccountAgent, |
| 509 | } |
| 510 | |
| 511 | #[derive(Clone, Deserialize, Serialize)] |
| 512 | #[serde(rename_all = "camelCase")] |
| 513 | struct AgentThread { |
| 514 | id: String, |
| 515 | agent_id: String, |
| 516 | #[serde(default)] |
| 517 | project_id: String, |
| 518 | title: String, |
| 519 | model: String, |
| 520 | model_provider: String, |
| 521 | #[serde(default)] |
| 522 | model_provider_id: String, |
| 523 | #[serde(default)] |
| 524 | kind: String, |
| 525 | #[serde(default)] |
| 526 | archived_at: String, |
| 527 | } |
| 528 | |
| 529 | #[derive(Deserialize)] |
| 530 | struct ThreadResponse { |
| 531 | thread: AgentThread, |
| 532 | } |
| 533 | |
| 534 | #[derive(Deserialize)] |
| 535 | struct TurnResponse { |
| 536 | turn: TurnReceipt, |
| 537 | } |
| 538 | |
| 539 | #[derive(Deserialize)] |
| 540 | struct TurnReceipt { |
| 541 | id: String, |
| 542 | status: String, |
| 543 | } |
| 544 | |
| 545 | #[derive(Deserialize)] |
| 546 | #[serde(rename_all = "camelCase")] |
| 547 | struct AccountAgentWork { |
| 548 | id: String, |
| 549 | agent_id: String, |
| 550 | status: String, |
| 551 | #[serde(default)] |
| 552 | objective: String, |
| 553 | } |
| 554 | |
| 555 | #[derive(Deserialize)] |
| 556 | #[serde(rename_all = "camelCase")] |
| 557 | struct AgentWorkMessageResponse { |
| 558 | intent: String, |
| 559 | /// Why the Agent read the message that way (a server-owned code). |
| 560 | #[serde(default)] |
| 561 | reason: String, |
| 562 | /// The action applied to active Work when `intent` is `control`. |
| 563 | #[serde(default)] |
| 564 | control_action: Option<String>, |
| 565 | #[serde(default)] |
| 566 | confident: Option<bool>, |
| 567 | /// What the Agent offers to do when it declined to act on an unclear message. |
| 568 | #[serde(default)] |
| 569 | suggestion: Option<String>, |
| 570 | work: Option<AccountAgentWork>, |
| 571 | #[serde(default)] |
| 572 | queued_work: Vec<AccountAgentWork>, |
| 573 | } |
| 574 | |
| 575 | #[derive(Deserialize)] |
| 576 | #[serde(rename_all = "camelCase")] |
| 577 | struct AgentTurnEvent { |
| 578 | seq: u64, |
| 579 | #[serde(rename = "type")] |
| 580 | kind: String, |
| 581 | turn_id: String, |
| 582 | payload: serde_json::Value, |
| 583 | } |
| 584 | |
| 585 | struct AgentTurnResult { |
| 586 | last_seq: u64, |
| 587 | status: String, |
| 588 | answer: String, |
| 589 | seen_turn: bool, |
| 590 | } |
| 591 | |
| 592 | impl CatalogProvider { |
| 593 | /// Non-empty display label, falling back to the id. |
| 594 | fn display_label(&self) -> String { |
| 595 | let label = printable(&self.label); |
| 596 | if label.is_empty() { |
| 597 | printable(&self.id) |
| 598 | } else { |
| 599 | label |
| 600 | } |
| 601 | } |
| 602 | |
| 603 | /// The local [`ProviderKind`] this catalog row maps onto, if any. |
| 604 | /// |
| 605 | /// The catalog states its own runtime mapping (`xiaomi` → |
| 606 | /// `xiaomi-mimo`); the row id is only a fallback for a provider whose |
| 607 | /// catalog id already equals the runtime id. |
| 608 | fn local_kind(&self) -> Option<ProviderKind> { |
| 609 | self.runtime_provider |
| 610 | .as_deref() |
| 611 | .map(str::trim) |
| 612 | .filter(|value| !value.is_empty()) |
| 613 | .and_then(ProviderKind::parse_config_identity) |
| 614 | .or_else(|| ProviderKind::parse_config_identity(&self.id)) |
| 615 | } |
| 616 | } |
| 617 | |
| 618 | /// Accept a provider id conservatively before it is ever put in a URL path. |
| 619 | /// |
| 620 | /// `^[a-z0-9][a-z0-9-]{0,63}$`. The catalog is remote data, so this guards |
| 621 | /// both directions: a hostile catalog cannot smuggle a path segment, and a |
| 622 | /// mistyped argument fails locally instead of as a confusing 404. |
| 623 | fn validate_provider_id(value: &str) -> Result<String> { |
| 624 | let trimmed = value.trim(); |
| 625 | let bytes = trimmed.as_bytes(); |
| 626 | let well_formed = (1..=64).contains(&bytes.len()) |
| 627 | && (bytes[0].is_ascii_lowercase() || bytes[0].is_ascii_digit()) |
| 628 | && bytes |
| 629 | .iter() |
| 630 | .all(|byte| byte.is_ascii_lowercase() || byte.is_ascii_digit() || *byte == b'-'); |
| 631 | if !well_formed { |
| 632 | bail!( |
| 633 | "`{}` is not a valid provider id. Ids are 1-64 characters of lowercase letters, digits, and `-`. Run `codewhale account keys list` to see the account's providers", |
| 634 | printable(trimmed) |
| 635 | ); |
| 636 | } |
| 637 | Ok(trimmed.to_string()) |
| 638 | } |
| 639 | |
| 640 | #[derive(Clone, Copy, Debug, PartialEq, Eq)] |
| 641 | enum HttpMethod { |
| 642 | Get, |
| 643 | Post, |
| 644 | Patch, |
| 645 | Put, |
| 646 | Delete, |
| 647 | } |
| 648 | |
| 649 | pub(crate) struct CloudRequest { |
| 650 | method: HttpMethod, |
| 651 | path: String, |
| 652 | bearer: Option<String>, |
| 653 | body: Option<Vec<u8>>, |
| 654 | } |
| 655 | |
| 656 | pub(crate) struct CloudResponse { |
| 657 | status: u16, |
| 658 | body: Vec<u8>, |
| 659 | /// `Retry-After` in whole seconds, when the service supplied one. Kept on |
| 660 | /// the response rather than re-parsed by callers so the retry policy has a |
| 661 | /// single source for how long the server asked us to wait. |
| 662 | retry_after: Option<u64>, |
| 663 | } |
| 664 | |
| 665 | pub(crate) trait CloudTransport { |
| 666 | fn execute(&self, request: CloudRequest) -> Result<CloudResponse>; |
| 667 | } |
| 668 | |
| 669 | struct ReqwestTransport { |
| 670 | base: Url, |
| 671 | client: reqwest::blocking::Client, |
| 672 | } |
| 673 | |
| 674 | impl ReqwestTransport { |
| 675 | fn new(base: Url) -> Result<Self> { |
| 676 | let client = codewhale_release::platform_blocking_http_client_builder() |
| 677 | .connect_timeout(Duration::from_secs(8)) |
| 678 | .timeout(Duration::from_secs(30)) |
| 679 | // Never replay bearer tokens or provider-key request bodies to a |
| 680 | // redirect target. The control-plane origin is an explicit trust |
| 681 | // boundary, so redirects are treated as ordinary non-2xx replies. |
| 682 | .redirect(reqwest::redirect::Policy::none()) |
| 683 | .user_agent(concat!("codewhale/", env!("CARGO_PKG_VERSION"))) |
| 684 | .build() |
| 685 | .context("failed to initialize the Codewhale account HTTP client")?; |
| 686 | Ok(Self { base, client }) |
| 687 | } |
| 688 | } |
| 689 | |
| 690 | impl CloudTransport for ReqwestTransport { |
| 691 | fn execute(&self, request: CloudRequest) -> Result<CloudResponse> { |
| 692 | let url = self |
| 693 | .base |
| 694 | .join(request.path.trim_start_matches('/')) |
| 695 | .context("failed to construct the Codewhale account request URL")?; |
| 696 | let method = match request.method { |
| 697 | HttpMethod::Get => reqwest::Method::GET, |
| 698 | HttpMethod::Post => reqwest::Method::POST, |
| 699 | HttpMethod::Patch => reqwest::Method::PATCH, |
| 700 | HttpMethod::Put => reqwest::Method::PUT, |
| 701 | HttpMethod::Delete => reqwest::Method::DELETE, |
| 702 | }; |
| 703 | let mut builder = self |
| 704 | .client |
| 705 | .request(method, url) |
| 706 | .header(reqwest::header::ACCEPT, "application/json"); |
| 707 | // Hosted launch waits for provider creation and guest bootstrap. Its |
| 708 | // HTTP deadline is separate from the guest's bounded trial runtime. |
| 709 | if request.method == HttpMethod::Post && request.path == "/api/cloud-sessions" { |
| 710 | builder = builder.timeout(Duration::from_secs(600)); |
| 711 | } |
| 712 | if let Some(token) = request.bearer { |
| 713 | builder = builder.bearer_auth(token); |
| 714 | } |
| 715 | if let Some(body) = request.body { |
| 716 | builder = builder |
| 717 | .header(reqwest::header::CONTENT_TYPE, "application/json") |
| 718 | .body(body); |
| 719 | } |
| 720 | let response = builder.send().map_err(|source| { |
| 721 | CloudTransportError::new("could not reach the Codewhale service", source) |
| 722 | })?; |
| 723 | let status = response.status().as_u16(); |
| 724 | let retry_after = response |
| 725 | .headers() |
| 726 | .get(reqwest::header::RETRY_AFTER) |
| 727 | .and_then(|value| value.to_str().ok()) |
| 728 | .and_then(|value| value.trim().parse::<u64>().ok()); |
| 729 | let mut body = Vec::new(); |
| 730 | response |
| 731 | .take(MAX_RESPONSE_BYTES + 1) |
| 732 | .read_to_end(&mut body) |
| 733 | .map_err(|source| { |
| 734 | CloudTransportError::new("failed to read the Codewhale service response", source) |
| 735 | })?; |
| 736 | if body.len() as u64 > MAX_RESPONSE_BYTES { |
| 737 | return Err(CloudTransportError::new( |
| 738 | "The Codewhale service returned an unexpectedly large response", |
| 739 | std::io::Error::other("response exceeded the account API size limit"), |
| 740 | ) |
| 741 | .into()); |
| 742 | } |
| 743 | Ok(CloudResponse { |
| 744 | status, |
| 745 | body, |
| 746 | retry_after, |
| 747 | }) |
| 748 | } |
| 749 | } |
| 750 | |
| 751 | #[derive(Deserialize)] |
| 752 | #[serde(rename_all = "camelCase")] |
| 753 | struct DeviceStart { |
| 754 | device_code: String, |
| 755 | user_code: String, |
| 756 | verification_uri: String, |
| 757 | verification_uri_complete: String, |
| 758 | expires_in: u64, |
| 759 | interval: u64, |
| 760 | } |
| 761 | |
| 762 | #[derive(Deserialize)] |
| 763 | struct MeResponse { |
| 764 | user: CloudUser, |
| 765 | } |
| 766 | |
| 767 | #[derive(Serialize)] |
| 768 | #[serde(rename_all = "camelCase")] |
| 769 | struct DeviceTokenRequest<'a> { |
| 770 | device_code: &'a str, |
| 771 | } |
| 772 | |
| 773 | #[derive(Serialize)] |
| 774 | #[serde(rename_all = "camelCase")] |
| 775 | struct RefreshRequest<'a> { |
| 776 | refresh_token: &'a str, |
| 777 | } |
| 778 | |
| 779 | #[derive(Serialize)] |
| 780 | struct ModelKeyRequest<'a> { |
| 781 | key: &'a str, |
| 782 | label: &'a str, |
| 783 | } |
| 784 | |
| 785 | pub(crate) struct CloudClient<'a, T: CloudTransport> { |
| 786 | transport: &'a T, |
| 787 | account_store: AccountSessionStore, |
| 788 | } |
| 789 | |
| 790 | impl<'a, T: CloudTransport> CloudClient<'a, T> { |
| 791 | fn new(transport: &'a T, secrets: &'a Secrets, profile: &str, api_base: &'a str) -> Self { |
| 792 | Self { |
| 793 | transport, |
| 794 | account_store: AccountSessionStore::new(secrets.clone(), Some(profile), api_base), |
| 795 | } |
| 796 | } |
| 797 | |
| 798 | fn start_device(&self) -> Result<DeviceStart> { |
| 799 | let response = self.transport.execute(CloudRequest { |
| 800 | method: HttpMethod::Post, |
| 801 | path: "/api/cli/device/start".to_string(), |
| 802 | bearer: None, |
| 803 | body: Some(b"{}".to_vec()), |
| 804 | })?; |
| 805 | expect_json(response, &[200]) |
| 806 | } |
| 807 | |
| 808 | fn poll_device( |
| 809 | &self, |
| 810 | device: &DeviceStart, |
| 811 | timeout: Duration, |
| 812 | sleep: &mut dyn FnMut(Duration), |
| 813 | ) -> Result<AuthBundle> { |
| 814 | validate_device_code(&device.device_code)?; |
| 815 | let server_lifetime = |
| 816 | Duration::from_secs(device.expires_in.clamp(1, MAX_LOGIN_TIMEOUT_SECONDS)); |
| 817 | // The Codewhale account service answers HTTP 202 while the code is |
| 818 | // still pending, so the first response is already meaningful: poll |
| 819 | // immediately and sleep afterwards. It has no slow_down. |
| 820 | let bundle = codewhale_config::device_code::DeviceCodePoll::new( |
| 821 | timeout.min(server_lifetime), |
| 822 | "Codewhale account login timed out; run `codewhale login` to try again", |
| 823 | ) |
| 824 | .interval_seconds(Some(device.interval)) |
| 825 | .max_interval_seconds(10) |
| 826 | .run(sleep, || { |
| 827 | let response = self.transport.execute(CloudRequest { |
| 828 | method: HttpMethod::Post, |
| 829 | path: "/api/cli/device/token".to_string(), |
| 830 | bearer: None, |
| 831 | body: Some(json_body(&DeviceTokenRequest { |
| 832 | device_code: &device.device_code, |
| 833 | })?), |
| 834 | })?; |
| 835 | match response.status { |
| 836 | 200 => { |
| 837 | let bundle: AuthBundle = parse_json_body(&response.body)?; |
| 838 | validate_auth_bundle(&bundle)?; |
| 839 | Ok(DevicePollOutcome::Complete(bundle)) |
| 840 | } |
| 841 | 202 => Ok(DevicePollOutcome::Pending), |
| 842 | _ => Err(response_error(&response)), |
| 843 | } |
| 844 | })?; |
| 845 | self.save_auth(bundle.clone())?; |
| 846 | Ok(bundle) |
| 847 | } |
| 848 | |
| 849 | fn load_auth(&self) -> Result<Option<StoredCloudAuth>> { |
| 850 | self.account_store.load().context( |
| 851 | "the local Codewhale account session is unreadable; run `codewhale account logout` and sign in again", |
| 852 | ) |
| 853 | } |
| 854 | |
| 855 | fn save_auth(&self, bundle: AuthBundle) -> Result<()> { |
| 856 | self.account_store |
| 857 | .save(bundle) |
| 858 | .context("failed to save the Codewhale account session in the local secret store") |
| 859 | } |
| 860 | |
| 861 | fn me(&self) -> Result<CloudUser> { |
| 862 | let (response, snapshot) = |
| 863 | self.execute_authenticated_snapshot(HttpMethod::Get, "/api/me", None)?; |
| 864 | let me: MeResponse = expect_json(response, &[200])?; |
| 865 | if me.user.id.trim().is_empty() { |
| 866 | bail!("The Codewhale service returned an account without an ID"); |
| 867 | } |
| 868 | if let Some(mut stored) = snapshot.load()? { |
| 869 | if stored |
| 870 | .bundle |
| 871 | .user |
| 872 | .as_ref() |
| 873 | .is_some_and(|user| !user.id.is_empty() && user.id != me.user.id) |
| 874 | { |
| 875 | bail!("The signed-in account changed. Refresh and try again."); |
| 876 | } |
| 877 | stored.bundle.user = Some(me.user.clone()); |
| 878 | if self |
| 879 | .account_store |
| 880 | .save_if_unchanged(&snapshot, stored.bundle)? |
| 881 | .is_none() |
| 882 | { |
| 883 | bail!("The signed-in account changed. Refresh and try again."); |
| 884 | } |
| 885 | } |
| 886 | Ok(me.user) |
| 887 | } |
| 888 | |
| 889 | fn set_key(&self, provider: &str, key: &str, label: &str) -> Result<()> { |
| 890 | let path = format!("/api/model-keys/{provider}"); |
| 891 | let response = self.execute_authenticated( |
| 892 | HttpMethod::Put, |
| 893 | &path, |
| 894 | Some(json_body(&ModelKeyRequest { key, label })?), |
| 895 | )?; |
| 896 | expect_empty(response, &[200, 201]) |
| 897 | } |
| 898 | |
| 899 | fn remove_key(&self, provider: &str) -> Result<()> { |
| 900 | let path = format!("/api/model-keys/{provider}"); |
| 901 | let response = self.execute_authenticated(HttpMethod::Delete, &path, None)?; |
| 902 | expect_empty(response, &[200, 204]) |
| 903 | } |
| 904 | |
| 905 | /// The account control plane's public provider catalog. |
| 906 | /// |
| 907 | /// This replaced a hardcoded eight-provider enum: the set of providers a |
| 908 | /// customer can connect is owned by the control plane, not the CLI, so a |
| 909 | /// newly supported provider must not need a CLI release. The route is |
| 910 | /// public, so no session is required to *list* what could be connected — |
| 911 | /// only to read or write this account's keys. |
| 912 | /// |
| 913 | /// Rows with an id this CLI would refuse to put in a URL path are dropped |
| 914 | /// rather than trusted; duplicates collapse onto the first row. |
| 915 | fn provider_catalog(&self) -> Result<Vec<CatalogProvider>> { |
| 916 | let response = self.transport.execute(CloudRequest { |
| 917 | method: HttpMethod::Get, |
| 918 | path: "/api/model-providers".to_string(), |
| 919 | bearer: None, |
| 920 | body: None, |
| 921 | })?; |
| 922 | let listing: ProviderCatalogResponse = expect_json(response, &[200])?; |
| 923 | let mut seen = std::collections::BTreeSet::new(); |
| 924 | let providers: Vec<CatalogProvider> = listing |
| 925 | .providers |
| 926 | .into_iter() |
| 927 | .filter(|row| validate_provider_id(&row.id).is_ok()) |
| 928 | .filter(|row| seen.insert(row.id.trim().to_string())) |
| 929 | .map(|mut row| { |
| 930 | row.id = row.id.trim().to_string(); |
| 931 | row |
| 932 | }) |
| 933 | .collect(); |
| 934 | if providers.is_empty() { |
| 935 | bail!( |
| 936 | "The Codewhale service returned no connectable providers. Check the account API origin, or try again" |
| 937 | ); |
| 938 | } |
| 939 | Ok(providers) |
| 940 | } |
| 941 | |
| 942 | fn computers(&self) -> Result<serde_json::Value> { |
| 943 | let response = self.execute_authenticated(HttpMethod::Get, "/api/computers", None)?; |
| 944 | expect_json(response, &[200]) |
| 945 | } |
| 946 | |
| 947 | fn agents(&self) -> Result<serde_json::Value> { |
| 948 | let response = self.execute_authenticated(HttpMethod::Get, "/api/agents", None)?; |
| 949 | expect_json(response, &[200]) |
| 950 | } |
| 951 | |
| 952 | fn projects(&self) -> Result<serde_json::Value> { |
| 953 | let response = self.execute_authenticated(HttpMethod::Get, "/api/projects", None)?; |
| 954 | expect_json(response, &[200]) |
| 955 | } |
| 956 | |
| 957 | fn github_bindings(&self) -> Result<serde_json::Value> { |
| 958 | let response = |
| 959 | self.execute_authenticated(HttpMethod::Get, "/api/integrations/github/bindings", None)?; |
| 960 | if matches!(response.status, 404 | 503) { |
| 961 | return Err(response_error(&response)).context( |
| 962 | "GitHub repository bindings are unavailable on this Codewhale API; the account GitHub route must be attached before CLI setup", |
| 963 | ); |
| 964 | } |
| 965 | expect_json(response, &[200]) |
| 966 | } |
| 967 | |
| 968 | fn create_github_project( |
| 969 | &self, |
| 970 | name: &str, |
| 971 | binding: &AccountGitHubBinding, |
| 972 | operation_key: &str, |
| 973 | ) -> Result<AccountProject> { |
| 974 | let name = validate_named_text(name, "Project name", 80)?; |
| 975 | let operation_key = validate_operation_key(operation_key)?; |
| 976 | let response = self.execute_authenticated( |
| 977 | HttpMethod::Post, |
| 978 | "/api/projects", |
| 979 | Some(json_body(&serde_json::json!({ |
| 980 | "name": name, |
| 981 | "defaultMode": "chat", |
| 982 | "chatFilesystem": "optional_scratch", |
| 983 | "repoBindingId": binding.id, |
| 984 | "operationKey": operation_key, |
| 985 | }))?), |
| 986 | )?; |
| 987 | let result: ProjectResponse = expect_json(response, &[200, 201])?; |
| 988 | validate_resource_id(&result.project.id, "Project")?; |
| 989 | if result.project.default_repo_provider != "github" |
| 990 | || !result |
| 991 | .project |
| 992 | .default_repo |
| 993 | .eq_ignore_ascii_case(&binding.repo) |
| 994 | { |
| 995 | bail!("The Codewhale service returned a Project for a different GitHub repository"); |
| 996 | } |
| 997 | Ok(result.project) |
| 998 | } |
| 999 | |
| 1000 | fn create_agent( |
| 1001 | &self, |
| 1002 | name: &str, |
| 1003 | project_id: Option<&str>, |
| 1004 | operation_key: &str, |
| 1005 | ) -> Result<AccountAgent> { |
| 1006 | let name = validate_named_text(name, "Agent name", 80)?; |
| 1007 | let operation_key = validate_operation_key(operation_key)?; |
| 1008 | let mut body = serde_json::json!({ |
| 1009 | "name": name, |
| 1010 | "operationKey": operation_key, |
| 1011 | }); |
| 1012 | if let Some(project_id) = project_id { |
| 1013 | body["projectId"] = serde_json::json!(validate_resource_id(project_id, "Project")?); |
| 1014 | } |
| 1015 | let response = |
| 1016 | self.execute_authenticated(HttpMethod::Post, "/api/agents", Some(json_body(&body)?))?; |
| 1017 | let result: AgentResponse = expect_json(response, &[200, 201])?; |
| 1018 | Ok(result.agent) |
| 1019 | } |
| 1020 | |
| 1021 | fn bind_agent_project( |
| 1022 | &self, |
| 1023 | agent_id: &str, |
| 1024 | project_id: &str, |
| 1025 | revision: u64, |
| 1026 | ) -> Result<AccountAgent> { |
| 1027 | let agent_id = validate_resource_id(agent_id, "Agent")?; |
| 1028 | let project_id = validate_resource_id(project_id, "Project")?; |
| 1029 | if revision == 0 { |
| 1030 | bail!("The Codewhale service returned an Agent without a valid revision"); |
| 1031 | } |
| 1032 | let response = self.execute_authenticated( |
| 1033 | HttpMethod::Patch, |
| 1034 | &format!("/api/agents/{agent_id}"), |
| 1035 | Some(json_body(&serde_json::json!({ |
| 1036 | "projectId": project_id, |
| 1037 | "revision": revision, |
| 1038 | }))?), |
| 1039 | )?; |
| 1040 | let result: AgentResponse = expect_json(response, &[200])?; |
| 1041 | if result.agent.id != agent_id || result.agent.project_id != project_id { |
| 1042 | bail!("The Codewhale service returned an unexpected Agent Project binding"); |
| 1043 | } |
| 1044 | Ok(result.agent) |
| 1045 | } |
| 1046 | |
| 1047 | fn threads(&self, agent_id: &str) -> Result<Vec<AgentThread>> { |
| 1048 | let agent_id = validate_resource_id(agent_id, "Agent")?; |
| 1049 | let response = self.execute_authenticated( |
| 1050 | HttpMethod::Get, |
| 1051 | &format!("/v1/threads/summary?agentId={agent_id}&limit=100"), |
| 1052 | None, |
| 1053 | )?; |
| 1054 | expect_json(response, &[200]) |
| 1055 | } |
| 1056 | |
| 1057 | fn thread(&self, id: &str) -> Result<AgentThread> { |
| 1058 | let id = validate_resource_id(id, "Conversation")?; |
| 1059 | let response = |
| 1060 | self.execute_authenticated(HttpMethod::Get, &format!("/v1/threads/{id}"), None)?; |
| 1061 | let result: ThreadResponse = expect_json(response, &[200])?; |
| 1062 | Ok(result.thread) |
| 1063 | } |
| 1064 | |
| 1065 | fn create_agent_thread( |
| 1066 | &self, |
| 1067 | agent_id: &str, |
| 1068 | project_id: Option<&str>, |
| 1069 | title: &str, |
| 1070 | provider: &str, |
| 1071 | model: &str, |
| 1072 | operation_key: &str, |
| 1073 | ) -> Result<AgentThread> { |
| 1074 | let agent_id = validate_resource_id(agent_id, "Agent")?; |
| 1075 | let title = validate_named_text(title, "Conversation title", 120)?; |
| 1076 | let (provider, model) = validate_model_route(provider, model)?; |
| 1077 | let operation_key = validate_operation_key(operation_key)?; |
| 1078 | let mut body = serde_json::json!({ |
| 1079 | "title": title, |
| 1080 | "productMode": "chat", |
| 1081 | "mode": "chat", |
| 1082 | "agentId": agent_id, |
| 1083 | "modelProvider": provider, |
| 1084 | "model": model, |
| 1085 | "operationKey": operation_key, |
| 1086 | }); |
| 1087 | if let Some(project_id) = project_id { |
| 1088 | body["projectId"] = serde_json::json!(validate_resource_id(project_id, "Project")?); |
| 1089 | } |
| 1090 | let response = |
| 1091 | self.execute_authenticated(HttpMethod::Post, "/v1/threads", Some(json_body(&body)?))?; |
| 1092 | let result: ThreadResponse = expect_json(response, &[200, 201])?; |
| 1093 | Ok(result.thread) |
| 1094 | } |
| 1095 | |
| 1096 | fn send_agent_turn( |
| 1097 | &self, |
| 1098 | thread: &AgentThread, |
| 1099 | prompt: &str, |
| 1100 | billing_mode: &str, |
| 1101 | operation_key: &str, |
| 1102 | ) -> Result<TurnReceipt> { |
| 1103 | let thread_id = validate_resource_id(&thread.id, "Conversation")?; |
| 1104 | let prompt = prompt.trim(); |
| 1105 | if prompt.is_empty() || prompt.chars().count() > 32_000 { |
| 1106 | bail!("Message must contain 1-32000 characters"); |
| 1107 | } |
| 1108 | let billing_mode = validate_billing_mode(billing_mode)?; |
| 1109 | let (provider, model) = validate_model_route(&thread.model_provider, &thread.model)?; |
| 1110 | let operation_key = validate_operation_key(operation_key)?; |
| 1111 | let response = self.execute_authenticated( |
| 1112 | HttpMethod::Post, |
| 1113 | &format!("/v1/threads/{thread_id}/turns"), |
| 1114 | Some(json_body(&serde_json::json!({ |
| 1115 | "prompt": prompt, |
| 1116 | "billingMode": billing_mode, |
| 1117 | "modelProvider": provider, |
| 1118 | "modelProviderId": thread.model_provider_id.as_str(), |
| 1119 | "model": model, |
| 1120 | "requiresByok": billing_mode == "byok_external", |
| 1121 | "mode": "chat", |
| 1122 | "productMode": "chat", |
| 1123 | "operationKey": operation_key, |
| 1124 | "sourceMessageId": operation_key, |
| 1125 | }))?), |
| 1126 | )?; |
| 1127 | let result: TurnResponse = expect_json(response, &[200, 201, 202])?; |
| 1128 | Ok(result.turn) |
| 1129 | } |
| 1130 | |
| 1131 | fn agent_turn_result( |
| 1132 | &self, |
| 1133 | thread_id: &str, |
| 1134 | turn_id: &str, |
| 1135 | since_seq: u64, |
| 1136 | ) -> Result<AgentTurnResult> { |
| 1137 | let thread_id = validate_resource_id(thread_id, "Conversation")?; |
| 1138 | let turn_id = validate_turn_id(turn_id)?; |
| 1139 | if since_seq > i64::MAX as u64 { |
| 1140 | bail!("Event sequence is too large"); |
| 1141 | } |
| 1142 | let response = self.execute_authenticated( |
| 1143 | HttpMethod::Get, |
| 1144 | &format!("/v1/threads/{thread_id}/events?since_seq={since_seq}"), |
| 1145 | None, |
| 1146 | )?; |
| 1147 | if response.status != 200 { |
| 1148 | return Err(response_error(&response)); |
| 1149 | } |
| 1150 | parse_agent_turn_events(&response.body, turn_id, since_seq) |
| 1151 | } |
| 1152 | |
| 1153 | fn assign_agent_work( |
| 1154 | &self, |
| 1155 | agent_id: &str, |
| 1156 | objective: &str, |
| 1157 | message_id: &str, |
| 1158 | ) -> Result<AgentWorkMessageResponse> { |
| 1159 | let agent_id = validate_resource_id(agent_id, "Agent")?; |
| 1160 | let objective = objective.trim(); |
| 1161 | if objective.is_empty() || objective.chars().count() > 32_000 { |
| 1162 | bail!("Work objective must contain 1-32000 characters"); |
| 1163 | } |
| 1164 | let message_id = validate_operation_key(message_id)?; |
| 1165 | let response = self.execute_authenticated( |
| 1166 | HttpMethod::Post, |
| 1167 | &format!("/api/agents/{agent_id}/messages"), |
| 1168 | Some(json_body(&serde_json::json!({ |
| 1169 | "messageId": message_id, |
| 1170 | "text": objective, |
| 1171 | }))?), |
| 1172 | )?; |
| 1173 | if ![200, 201, 202].contains(&response.status) { |
| 1174 | return Err(response_error(&response)); |
| 1175 | } |
| 1176 | // A 2xx the client cannot read means the service acted: an unknown |
| 1177 | // outcome (replay with the same message id), not a plain failure. |
| 1178 | serde_json::from_slice(&response.body).map_err(|source| { |
| 1179 | CloudTransportError::new( |
| 1180 | "The Codewhale service returned an unreadable reply to this message", |
| 1181 | source, |
| 1182 | ) |
| 1183 | .into() |
| 1184 | }) |
| 1185 | } |
| 1186 | |
| 1187 | fn agent_work_status(&self, id: &str) -> Result<serde_json::Value> { |
| 1188 | let id = validate_resource_id(id, "Work")?; |
| 1189 | let response = |
| 1190 | self.execute_authenticated(HttpMethod::Get, &format!("/api/runs/{id}"), None)?; |
| 1191 | expect_json(response, &[200]) |
| 1192 | } |
| 1193 | |
| 1194 | fn computer(&self, id: &str) -> Result<serde_json::Value> { |
| 1195 | let id = validate_computer_id(id)?; |
| 1196 | let response = |
| 1197 | self.execute_authenticated(HttpMethod::Get, &format!("/api/computers/{id}"), None)?; |
| 1198 | expect_json(response, &[200]) |
| 1199 | } |
| 1200 | |
| 1201 | fn computer_usage(&self, id: &str) -> Result<serde_json::Value> { |
| 1202 | let id = validate_computer_id(id)?; |
| 1203 | let response = self.execute_authenticated( |
| 1204 | HttpMethod::Get, |
| 1205 | &format!("/api/computers/{id}/usage"), |
| 1206 | None, |
| 1207 | )?; |
| 1208 | expect_json(response, &[200]) |
| 1209 | } |
| 1210 | |
| 1211 | fn create_computer( |
| 1212 | &self, |
| 1213 | name: &str, |
| 1214 | boat_trial: bool, |
| 1215 | eu_compute_opt_in: bool, |
| 1216 | ) -> Result<AccountComputer> { |
| 1217 | if boat_trial != eu_compute_opt_in { |
| 1218 | bail!("Boat trial requires both --boat-trial and --eu-compute-opt-in"); |
| 1219 | } |
| 1220 | let name = validate_computer_name(name)?; |
| 1221 | let response = self.execute_authenticated( |
| 1222 | HttpMethod::Post, |
| 1223 | "/api/computers", |
| 1224 | Some(json_body(&ComputerCreateRequest { |
| 1225 | name: &name, |
| 1226 | provider: boat_trial.then_some("boat"), |
| 1227 | boat_eu_compute_opt_in: boat_trial.then_some(true), |
| 1228 | })?), |
| 1229 | )?; |
| 1230 | let result: ComputerResponse = expect_json(response, &[200, 201])?; |
| 1231 | Ok(result.computer) |
| 1232 | } |
| 1233 | |
| 1234 | fn computer_action(&self, id: &str, action: &str) -> Result<ComputerResponse> { |
| 1235 | let id = validate_computer_id(id)?; |
| 1236 | let response = self.execute_authenticated( |
| 1237 | HttpMethod::Post, |
| 1238 | &format!("/api/computers/{id}/{action}"), |
| 1239 | None, |
| 1240 | )?; |
| 1241 | let status = response.status; |
| 1242 | let result: ComputerResponse = expect_json( |
| 1243 | response, |
| 1244 | if action == "start" { |
| 1245 | &[200, 202] |
| 1246 | } else { |
| 1247 | &[200] |
| 1248 | }, |
| 1249 | )?; |
| 1250 | if status == 202 && !result.queued { |
| 1251 | bail!("The Codewhale service returned a queued start without queue details"); |
| 1252 | } |
| 1253 | Ok(result) |
| 1254 | } |
| 1255 | |
| 1256 | fn delete_computer(&self, id: &str) -> Result<()> { |
| 1257 | let id = validate_computer_id(id)?; |
| 1258 | let response = |
| 1259 | self.execute_authenticated(HttpMethod::Delete, &format!("/api/computers/{id}"), None)?; |
| 1260 | let result: ComputerDeleteResponse = expect_json(response, &[200])?; |
| 1261 | if !result.deleted || result.computer_id != id { |
| 1262 | bail!("The Codewhale service did not confirm deletion of Computer {id}"); |
| 1263 | } |
| 1264 | Ok(()) |
| 1265 | } |
| 1266 | |
| 1267 | fn logout(&self) -> Result<bool> { |
| 1268 | let snapshot = self.account_store.snapshot()?; |
| 1269 | self.account_store |
| 1270 | .with_transaction(|transaction| -> Result<bool> { |
| 1271 | if !transaction.matches(&snapshot) { |
| 1272 | bail!("The signed-in account changed. Refresh and try again."); |
| 1273 | } |
| 1274 | let stored = match transaction.load() { |
| 1275 | Ok(Some(stored)) => stored, |
| 1276 | Ok(None) | Err(_) => { |
| 1277 | transaction.clear(); |
| 1278 | return Ok(false); |
| 1279 | } |
| 1280 | }; |
| 1281 | let body = json_body(&RefreshRequest { |
| 1282 | refresh_token: &stored.bundle.refresh_token, |
| 1283 | })?; |
| 1284 | let response = self.transport.execute(CloudRequest { |
| 1285 | method: HttpMethod::Post, |
| 1286 | path: "/api/auth/logout".into(), |
| 1287 | bearer: None, |
| 1288 | body: Some(body), |
| 1289 | })?; |
| 1290 | if (200..300).contains(&response.status) || matches!(response.status, 401 | 403) { |
| 1291 | transaction.clear(); |
| 1292 | return Ok((200..300).contains(&response.status)); |
| 1293 | } |
| 1294 | // Keep custody on transient failure so revocation can be retried. |
| 1295 | Err(response_error(&response)) |
| 1296 | }) |
| 1297 | } |
| 1298 | |
| 1299 | fn execute_authenticated( |
| 1300 | &self, |
| 1301 | method: HttpMethod, |
| 1302 | path: &str, |
| 1303 | body: Option<Vec<u8>>, |
| 1304 | ) -> Result<CloudResponse> { |
| 1305 | self.execute_authenticated_snapshot(method, path, body) |
| 1306 | .map(|(response, _)| response) |
| 1307 | } |
| 1308 | |
| 1309 | fn execute_authenticated_snapshot( |
| 1310 | &self, |
| 1311 | method: HttpMethod, |
| 1312 | path: &str, |
| 1313 | body: Option<Vec<u8>>, |
| 1314 | ) -> Result<(CloudResponse, AccountSessionSnapshot)> { |
| 1315 | let snapshot = self.account_store.snapshot()?; |
| 1316 | let Some(mut stored) = snapshot.load()? else { |
| 1317 | bail!("Not signed in. Run `codewhale login` first"); |
| 1318 | }; |
| 1319 | let first = self.transport.execute(CloudRequest { |
| 1320 | method, |
| 1321 | path: path.into(), |
| 1322 | bearer: Some(stored.bundle.access_token.clone()), |
| 1323 | body: body.clone(), |
| 1324 | })?; |
| 1325 | if first.status != 401 { |
| 1326 | return Ok((first, snapshot)); |
| 1327 | } |
| 1328 | // Serialize the refresh HTTP request itself with native/CLI writers: |
| 1329 | // two processes must not spend the same rotating refresh token. |
| 1330 | let renewed = self |
| 1331 | .account_store |
| 1332 | .with_transaction(|transaction| -> Result<_> { |
| 1333 | if !transaction.matches(&snapshot) { |
| 1334 | bail!("The signed-in account changed. Refresh and try again."); |
| 1335 | } |
| 1336 | let refresh = self.transport.execute(CloudRequest { |
| 1337 | method: HttpMethod::Post, |
| 1338 | path: "/api/auth/refresh".into(), |
| 1339 | bearer: None, |
| 1340 | body: Some(json_body(&RefreshRequest { |
| 1341 | refresh_token: &stored.bundle.refresh_token, |
| 1342 | })?), |
| 1343 | })?; |
| 1344 | match refresh.status { |
| 1345 | 200 => {} |
| 1346 | 401 => { |
| 1347 | transaction.clear(); |
| 1348 | return Ok(None); |
| 1349 | } |
| 1350 | _ => return Err(response_error(&refresh)), |
| 1351 | } |
| 1352 | let mut next: AuthBundle = parse_json_body(&refresh.body)?; |
| 1353 | validate_auth_bundle(&next)?; |
| 1354 | if next.user.is_none() { |
| 1355 | next.user = stored.bundle.user.take(); |
| 1356 | } |
| 1357 | if next.session.is_none() { |
| 1358 | next.session = stored.bundle.session.take(); |
| 1359 | } |
| 1360 | transaction.replace(next.clone())?; |
| 1361 | Ok(Some((next, transaction.snapshot()))) |
| 1362 | })?; |
| 1363 | let Some((next, next_snapshot)) = renewed else { |
| 1364 | bail!("The Codewhale account session expired. Run `codewhale login` again"); |
| 1365 | }; |
| 1366 | // The rotated token is durable before a potentially failing retry. |
| 1367 | let retried = self.transport.execute(CloudRequest { |
| 1368 | method, |
| 1369 | path: path.into(), |
| 1370 | bearer: Some(next.access_token), |
| 1371 | body, |
| 1372 | })?; |
| 1373 | if retried.status == 401 { |
| 1374 | self.account_store.clear_if_unchanged(&next_snapshot)?; |
| 1375 | bail!("The Codewhale account session expired. Run `codewhale login` again"); |
| 1376 | } |
| 1377 | Ok((retried, next_snapshot)) |
| 1378 | } |
| 1379 | |
| 1380 | /// Whether an interactive session exists for this profile and origin. |
| 1381 | /// |
| 1382 | /// A management command asks this before it asks anything of the network, |
| 1383 | /// so "you have a machine key but no login" is answered locally instead of |
| 1384 | /// as a 403 from a route the key was never allowed to touch. |
| 1385 | fn has_session(&self) -> Result<bool> { |
| 1386 | Ok(self.load_auth()?.is_some()) |
| 1387 | } |
| 1388 | |
| 1389 | /// `execute_authenticated`, retrying only what the caller marks replayable. |
| 1390 | /// |
| 1391 | /// `machine::Retry::Never` is not a default worth having: the one POST in |
| 1392 | /// this surface mints a secret shown exactly once, so a retry that quietly |
| 1393 | /// succeeded server-side would leave an unrevocable key behind. |
| 1394 | fn execute_authenticated_with_retry( |
| 1395 | &self, |
| 1396 | method: HttpMethod, |
| 1397 | path: &str, |
| 1398 | body: Option<Vec<u8>>, |
| 1399 | retry: machine::Retry, |
| 1400 | sleeper: &mut dyn FnMut(Duration), |
| 1401 | ) -> Result<CloudResponse> { |
| 1402 | let max_attempts = if retry == machine::Retry::Idempotent { |
| 1403 | 3 |
| 1404 | } else { |
| 1405 | 1 |
| 1406 | }; |
| 1407 | let mut attempt = 1; |
| 1408 | loop { |
| 1409 | let response = self.execute_authenticated(method, path, body.clone())?; |
| 1410 | if (200..300).contains(&response.status) || attempt >= max_attempts { |
| 1411 | return Ok(response); |
| 1412 | } |
| 1413 | let retry_after = response.retry_after; |
| 1414 | if !machine::classify(&response).retryable { |
| 1415 | return Ok(response); |
| 1416 | } |
| 1417 | sleeper(machine::backoff_delay(attempt, retry_after)); |
| 1418 | attempt += 1; |
| 1419 | } |
| 1420 | } |
| 1421 | } |
| 1422 | |
| 1423 | fn validate_computer_id(value: &str) -> Result<&str> { |
| 1424 | validate_resource_id(value, "Computer") |
| 1425 | } |
| 1426 | |
| 1427 | fn validate_resource_id<'a>(value: &'a str, kind: &str) -> Result<&'a str> { |
| 1428 | let id = value.trim(); |
| 1429 | let bytes = id.as_bytes(); |
| 1430 | if bytes.is_empty() |
| 1431 | || bytes.len() > 160 |
| 1432 | || !bytes[0].is_ascii_alphanumeric() |
| 1433 | || !bytes |
| 1434 | .iter() |
| 1435 | .all(|byte| byte.is_ascii_alphanumeric() || matches!(byte, b'-' | b'_')) |
| 1436 | { |
| 1437 | bail!("{kind} ID must be a bounded identifier of letters, digits, `-`, or `_`"); |
| 1438 | } |
| 1439 | Ok(id) |
| 1440 | } |
| 1441 | |
| 1442 | fn validate_turn_id(value: &str) -> Result<&str> { |
| 1443 | let id = value.trim(); |
| 1444 | let bytes = id.as_bytes(); |
| 1445 | if bytes.is_empty() |
| 1446 | || bytes.len() > 160 |
| 1447 | || !bytes[0].is_ascii_alphanumeric() |
| 1448 | || !bytes.iter().all(|byte| { |
| 1449 | byte.is_ascii_alphanumeric() || matches!(byte, b'-' | b'_' | b'.' | b':' | b'@') |
| 1450 | }) |
| 1451 | || id.contains("..") |
| 1452 | { |
| 1453 | bail!( |
| 1454 | "Turn ID must be a bounded identifier of letters, digits, `-`, `_`, `.`, `:`, or `@`" |
| 1455 | ); |
| 1456 | } |
| 1457 | Ok(id) |
| 1458 | } |
| 1459 | |
| 1460 | fn validate_computer_name(value: &str) -> Result<String> { |
| 1461 | validate_named_text(value, "Computer name", 80) |
| 1462 | } |
| 1463 | |
| 1464 | fn validate_named_text(value: &str, label: &str, maximum: usize) -> Result<String> { |
| 1465 | let text = value.trim(); |
| 1466 | if text.is_empty() || text.chars().count() > maximum || text.chars().any(char::is_control) { |
| 1467 | bail!("{label} must contain 1-{maximum} characters without control characters"); |
| 1468 | } |
| 1469 | Ok(text.to_string()) |
| 1470 | } |
| 1471 | |
| 1472 | fn validate_operation_key(value: &str) -> Result<&str> { |
| 1473 | let key = value.trim(); |
| 1474 | if key.is_empty() |
| 1475 | || key.len() > 128 |
| 1476 | || !key |
| 1477 | .bytes() |
| 1478 | .all(|byte| byte.is_ascii_alphanumeric() || matches!(byte, b'.' | b'_' | b':' | b'-')) |
| 1479 | { |
| 1480 | bail!("Operation key must be 1-128 URL-safe characters"); |
| 1481 | } |
| 1482 | Ok(key) |
| 1483 | } |
| 1484 | |
| 1485 | fn validate_billing_mode(value: &str) -> Result<&str> { |
| 1486 | match value.trim() { |
| 1487 | "byok_external" => Ok("byok_external"), |
| 1488 | "membership_included" | "managed_wallet" => { |
| 1489 | bail!( |
| 1490 | "Codewhale-managed model billing is unavailable under the current launch policy. Use --billing-mode byok_external with your own provider key" |
| 1491 | ) |
| 1492 | } |
| 1493 | _ => bail!("Billing mode must be byok_external under the current launch policy"), |
| 1494 | } |
| 1495 | } |
| 1496 | |
| 1497 | fn validate_model_route<'a>(provider: &'a str, model: &'a str) -> Result<(&'a str, &'a str)> { |
| 1498 | let provider = provider.trim(); |
| 1499 | let model = model.trim(); |
| 1500 | if provider.is_empty() |
| 1501 | || provider.len() > 128 |
| 1502 | || !provider |
| 1503 | .bytes() |
| 1504 | .all(|byte| byte.is_ascii_alphanumeric() || matches!(byte, b'.' | b'_' | b'-')) |
| 1505 | || model.is_empty() |
| 1506 | || model.len() > 256 |
| 1507 | || !model.bytes().all(|byte| { |
| 1508 | byte.is_ascii_alphanumeric() |
| 1509 | || matches!(byte, b'.' | b'_' | b':' | b'/' | b'@' | b'+' | b'-') |
| 1510 | }) |
| 1511 | || provider.contains("..") |
| 1512 | || model.contains("..") |
| 1513 | { |
| 1514 | bail!("Model route must contain a bounded provider and model identifier"); |
| 1515 | } |
| 1516 | Ok((provider, model)) |
| 1517 | } |
| 1518 | |
| 1519 | /// Refuse a (provider, model) the live account catalog does not list. |
| 1520 | /// |
| 1521 | /// The catalog is served data (`GET /api/model-providers`), so a new model |
| 1522 | /// needs no CLI release and a retired one stops being accepted here. The |
| 1523 | /// message names the catalog and what it does list for the provider. |
| 1524 | /// |
| 1525 | /// Returns the row's canonical id: the control plane stores that id even when |
| 1526 | /// the caller named the runtime alias (`xiaomi-mimo` for `xiaomi`), so it is |
| 1527 | /// the value to send and to compare the created conversation against. |
| 1528 | fn assert_catalog_route<'a>( |
| 1529 | catalog: &'a [CatalogProvider], |
| 1530 | provider: &str, |
| 1531 | model: &str, |
| 1532 | ) -> Result<&'a str> { |
| 1533 | let row = catalog.iter().find(|row| { |
| 1534 | row.connection_available != Some(false) |
| 1535 | && (row.id == provider || row.runtime_provider.as_deref() == Some(provider)) |
| 1536 | }); |
| 1537 | if let Some(row) = row |
| 1538 | && row.models.iter().any(|listed| listed == model) |
| 1539 | { |
| 1540 | return Ok(row.id.as_str()); |
| 1541 | } |
| 1542 | let listed = match row { |
| 1543 | Some(row) if !row.models.is_empty() => format!( |
| 1544 | "The catalog lists these models for `{}`: {}", |
| 1545 | printable(provider), |
| 1546 | row.models |
| 1547 | .iter() |
| 1548 | .take(12) |
| 1549 | .map(|model| printable(model)) |
| 1550 | .collect::<Vec<_>>() |
| 1551 | .join(", ") |
| 1552 | ), |
| 1553 | Some(_) => format!("The catalog lists no models for `{}`", printable(provider)), |
| 1554 | None => format!( |
| 1555 | "`{}` is not a hosted provider in the catalog. Known providers: {}", |
| 1556 | printable(provider), |
| 1557 | catalog |
| 1558 | .iter() |
| 1559 | .filter(|row| row.connection_available != Some(false)) |
| 1560 | .map(|row| row.id.as_str()) |
| 1561 | .collect::<Vec<_>>() |
| 1562 | .join(", ") |
| 1563 | ), |
| 1564 | }; |
| 1565 | bail!( |
| 1566 | "`{}/{}` is not listed in the live Codewhale model catalog (GET /api/model-providers), so no conversation was created. {listed}", |
| 1567 | printable(provider), |
| 1568 | printable(model) |
| 1569 | ) |
| 1570 | } |
| 1571 | |
| 1572 | fn resolve_account_agent(agents: &[AccountAgent], selector: &str) -> Result<AccountAgent> { |
| 1573 | let selector = selector.trim(); |
| 1574 | if selector.is_empty() { |
| 1575 | bail!("Choose an Agent name or ID"); |
| 1576 | } |
| 1577 | let by_id = agents.iter().find(|agent| agent.id == selector); |
| 1578 | let matches = if let Some(agent) = by_id { |
| 1579 | vec![agent] |
| 1580 | } else { |
| 1581 | agents |
| 1582 | .iter() |
| 1583 | .filter(|agent| agent.name == selector) |
| 1584 | .collect::<Vec<_>>() |
| 1585 | }; |
| 1586 | if matches.len() != 1 { |
| 1587 | bail!( |
| 1588 | "Expected one Agent named `{}`; found {}. Run `codewhale account agents list` and use its ID", |
| 1589 | printable(selector), |
| 1590 | matches.len() |
| 1591 | ); |
| 1592 | } |
| 1593 | let agent = matches[0]; |
| 1594 | validate_resource_id(&agent.id, "Agent")?; |
| 1595 | if agent.status != "active" { |
| 1596 | bail!("Agent {} is not active", printable(&agent.name)); |
| 1597 | } |
| 1598 | Ok(agent.clone()) |
| 1599 | } |
| 1600 | |
| 1601 | fn active_agent_thread(thread: &AgentThread, agent_id: &str) -> bool { |
| 1602 | thread.agent_id == agent_id && thread.archived_at.is_empty() && thread.kind != "smoke" |
| 1603 | } |
| 1604 | |
| 1605 | fn parse_agent_turn_events(body: &[u8], turn_id: &str, since_seq: u64) -> Result<AgentTurnResult> { |
| 1606 | let text = std::str::from_utf8(body) |
| 1607 | .context("The Codewhale service returned invalid conversation events")?; |
| 1608 | let normalized = text.replace("\r\n", "\n"); |
| 1609 | let mut result = AgentTurnResult { |
| 1610 | last_seq: since_seq, |
| 1611 | status: "not_seen".to_string(), |
| 1612 | answer: String::new(), |
| 1613 | seen_turn: false, |
| 1614 | }; |
| 1615 | for frame in normalized.split("\n\n") { |
| 1616 | let mut data = None; |
| 1617 | for line in frame.lines() { |
| 1618 | if let Some(value) = line.strip_prefix("data:") |
| 1619 | && data.replace(value.trim_start_matches(' ')).is_some() |
| 1620 | { |
| 1621 | bail!("The Codewhale service returned a malformed conversation event"); |
| 1622 | } |
| 1623 | } |
| 1624 | let Some(data) = data else { continue }; |
| 1625 | let event: AgentTurnEvent = serde_json::from_str(data) |
| 1626 | .context("The Codewhale service returned an invalid conversation event")?; |
| 1627 | if event.seq <= result.last_seq { |
| 1628 | bail!("The Codewhale service returned an out-of-order conversation event"); |
| 1629 | } |
| 1630 | result.last_seq = event.seq; |
| 1631 | if event.turn_id != turn_id { |
| 1632 | continue; |
| 1633 | } |
| 1634 | result.seen_turn = true; |
| 1635 | if result.status == "not_seen" { |
| 1636 | result.status = "pending".to_string(); |
| 1637 | } |
| 1638 | match event.kind.as_str() { |
| 1639 | "assistant.delta" => { |
| 1640 | let delta = event |
| 1641 | .payload |
| 1642 | .get("text") |
| 1643 | .and_then(serde_json::Value::as_str) |
| 1644 | .ok_or_else(|| { |
| 1645 | anyhow!("The Codewhale service returned invalid assistant text") |
| 1646 | })?; |
| 1647 | result.answer.push_str(delta); |
| 1648 | if result.answer.len() > 128 * 1024 { |
| 1649 | bail!("The Codewhale service returned an unexpectedly long answer"); |
| 1650 | } |
| 1651 | } |
| 1652 | "turn.completed" => { |
| 1653 | let status = event |
| 1654 | .payload |
| 1655 | .get("status") |
| 1656 | .and_then(serde_json::Value::as_str) |
| 1657 | .unwrap_or("completed"); |
| 1658 | if !matches!(status, "completed" | "failed") { |
| 1659 | bail!("The Codewhale service returned an invalid turn status"); |
| 1660 | } |
| 1661 | result.status = status.to_string(); |
| 1662 | } |
| 1663 | "turn.canceled" => result.status = "canceled".to_string(), |
| 1664 | _ => {} |
| 1665 | } |
| 1666 | } |
| 1667 | Ok(result) |
| 1668 | } |
| 1669 | |
| 1670 | fn write_agent_thread<W: Write>(out: &mut W, thread: &AgentThread) -> Result<()> { |
| 1671 | validate_resource_id(&thread.id, "Conversation")?; |
| 1672 | let (provider, model) = validate_model_route(&thread.model_provider, &thread.model)?; |
| 1673 | writeln!(out, "Conversation: {}", printable(&thread.title))?; |
| 1674 | writeln!(out, "ID: {}", thread.id)?; |
| 1675 | writeln!(out, "Model: {provider}/{model}")?; |
| 1676 | Ok(()) |
| 1677 | } |
| 1678 | |
| 1679 | fn run_projects<T: CloudTransport, W: Write>( |
| 1680 | command: CloudProjectsCommand, |
| 1681 | client: &CloudClient<'_, T>, |
| 1682 | machine: &machine::MachineKeyEnv, |
| 1683 | out: &mut W, |
| 1684 | ) -> Result<()> { |
| 1685 | if machine.is_present() { |
| 1686 | bail!( |
| 1687 | "Projects require an interactive Codewhale account login; unset CODEWHALE_API_KEY and run `codewhale login`" |
| 1688 | ); |
| 1689 | } |
| 1690 | match command { |
| 1691 | CloudProjectsCommand::List { json } => { |
| 1692 | let response = client.projects()?; |
| 1693 | if json { |
| 1694 | return write_computer_json(out, &response); |
| 1695 | } |
| 1696 | let listing: ProjectListResponse = serde_json::from_value(response) |
| 1697 | .context("The Codewhale service returned an invalid Project list")?; |
| 1698 | writeln!(out, "Codewhale Projects ({})", listing.projects.len())?; |
| 1699 | for project in listing.projects { |
| 1700 | validate_resource_id(&project.id, "Project")?; |
| 1701 | writeln!(out, "{} — {}", printable(&project.name), project.id)?; |
| 1702 | } |
| 1703 | Ok(()) |
| 1704 | } |
| 1705 | CloudProjectsCommand::Create { |
| 1706 | name, |
| 1707 | repo_binding_id, |
| 1708 | operation_key, |
| 1709 | } => { |
| 1710 | let listing: GitHubBindingListResponse = |
| 1711 | serde_json::from_value(client.github_bindings()?) |
| 1712 | .context("The Codewhale service returned an invalid GitHub repository list")?; |
| 1713 | let binding = listing |
| 1714 | .bindings |
| 1715 | .iter() |
| 1716 | .find(|binding| binding.id == repo_binding_id) |
| 1717 | .ok_or_else(|| anyhow!( |
| 1718 | "GitHub binding {} is not connected to this account. Run `codewhale account github bindings`", |
| 1719 | printable(&repo_binding_id) |
| 1720 | ))?; |
| 1721 | if binding.provider != "github" |
| 1722 | || binding.installation_id.is_empty() |
| 1723 | || ["error", "revoked", "suspended", "disabled"] |
| 1724 | .contains(&binding.status.to_ascii_lowercase().as_str()) |
| 1725 | { |
| 1726 | bail!( |
| 1727 | "That GitHub binding is unavailable for a Project. Reconnect the repository, then run `codewhale account github bindings`" |
| 1728 | ); |
| 1729 | } |
| 1730 | let project = client.create_github_project(&name, binding, &operation_key)?; |
| 1731 | writeln!(out, "Project: {}", printable(&project.name))?; |
| 1732 | writeln!(out, "ID: {}", project.id)?; |
| 1733 | writeln!( |
| 1734 | out, |
| 1735 | "GitHub repository: {}", |
| 1736 | printable(&project.default_repo) |
| 1737 | )?; |
| 1738 | writeln!( |
| 1739 | out, |
| 1740 | "Assign a Whale: codewhale account agents bind-project AGENT {}", |
| 1741 | project.id |
| 1742 | )?; |
| 1743 | Ok(()) |
| 1744 | } |
| 1745 | } |
| 1746 | } |
| 1747 | |
| 1748 | fn run_github<T: CloudTransport, W: Write>( |
| 1749 | command: CloudGithubCommand, |
| 1750 | client: &CloudClient<'_, T>, |
| 1751 | machine: &machine::MachineKeyEnv, |
| 1752 | out: &mut W, |
| 1753 | ) -> Result<()> { |
| 1754 | if machine.is_present() { |
| 1755 | bail!( |
| 1756 | "GitHub bindings require an interactive Codewhale account login; unset CODEWHALE_API_KEY and run `codewhale login`" |
| 1757 | ); |
| 1758 | } |
| 1759 | match command { |
| 1760 | CloudGithubCommand::Bindings { json } => { |
| 1761 | let response = client.github_bindings()?; |
| 1762 | if json { |
| 1763 | return write_computer_json(out, &response); |
| 1764 | } |
| 1765 | let listing: GitHubBindingListResponse = serde_json::from_value(response) |
| 1766 | .context("The Codewhale service returned an invalid GitHub repository list")?; |
| 1767 | writeln!(out, "GitHub repositories ({})", listing.bindings.len())?; |
| 1768 | if listing.bindings.is_empty() { |
| 1769 | writeln!( |
| 1770 | out, |
| 1771 | "No GitHub repositories are connected to this account. Authorize one in Codewhale's GitHub setup, then retry." |
| 1772 | )?; |
| 1773 | } |
| 1774 | for binding in listing.bindings { |
| 1775 | if binding.provider != "github" || binding.id.is_empty() || binding.repo.is_empty() |
| 1776 | { |
| 1777 | bail!("The Codewhale service returned an invalid GitHub repository binding"); |
| 1778 | } |
| 1779 | writeln!( |
| 1780 | out, |
| 1781 | "{} — {} ({})", |
| 1782 | printable(&binding.repo), |
| 1783 | printable(&binding.id), |
| 1784 | printable(&binding.status) |
| 1785 | )?; |
| 1786 | } |
| 1787 | Ok(()) |
| 1788 | } |
| 1789 | CloudGithubCommand::Bind { |
| 1790 | repo, |
| 1791 | installation_id, |
| 1792 | } => work::bind_github_repo(client, out, &repo, installation_id.as_deref()), |
| 1793 | } |
| 1794 | } |
| 1795 | |
| 1796 | fn run_agents<T: CloudTransport, W: Write>( |
| 1797 | command: CloudAgentsCommand, |
| 1798 | client: &CloudClient<'_, T>, |
| 1799 | machine: &machine::MachineKeyEnv, |
| 1800 | out: &mut W, |
| 1801 | ) -> Result<()> { |
| 1802 | if machine.is_present() { |
| 1803 | bail!( |
| 1804 | "Agents require an interactive Codewhale account login; unset CODEWHALE_API_KEY and run `codewhale login`" |
| 1805 | ); |
| 1806 | } |
| 1807 | match command { |
| 1808 | CloudAgentsCommand::List { json } => { |
| 1809 | let response = client.agents()?; |
| 1810 | if json { |
| 1811 | return write_computer_json(out, &response); |
| 1812 | } |
| 1813 | let listing: AgentListResponse = serde_json::from_value(response) |
| 1814 | .context("The Codewhale service returned an invalid Agent list")?; |
| 1815 | writeln!(out, "Codewhale Agents ({})", listing.agents.len())?; |
| 1816 | for agent in listing.agents { |
| 1817 | validate_resource_id(&agent.id, "Agent")?; |
| 1818 | writeln!(out, "{} — {}", printable(&agent.name), agent.id)?; |
| 1819 | } |
| 1820 | Ok(()) |
| 1821 | } |
| 1822 | CloudAgentsCommand::Create { |
| 1823 | name, |
| 1824 | project_id, |
| 1825 | operation_key, |
| 1826 | } => { |
| 1827 | let agent = client.create_agent(&name, project_id.as_deref(), &operation_key)?; |
| 1828 | validate_resource_id(&agent.id, "Agent")?; |
| 1829 | if let Some(project_id) = project_id.as_deref() |
| 1830 | && agent.project_id != validate_resource_id(project_id, "Project")? |
| 1831 | { |
| 1832 | bail!("The Codewhale service returned an unexpected Agent Project binding"); |
| 1833 | } |
| 1834 | writeln!(out, "Agent: {}", printable(&agent.name))?; |
| 1835 | writeln!(out, "ID: {}", agent.id)?; |
| 1836 | if !agent.project_id.is_empty() { |
| 1837 | writeln!( |
| 1838 | out, |
| 1839 | "Project ID: {}", |
| 1840 | validate_resource_id(&agent.project_id, "Project")? |
| 1841 | )?; |
| 1842 | } |
| 1843 | writeln!( |
| 1844 | out, |
| 1845 | "Create request ID: {}", |
| 1846 | validate_operation_key(&operation_key)? |
| 1847 | )?; |
| 1848 | if agent.project_id.is_empty() { |
| 1849 | writeln!( |
| 1850 | out, |
| 1851 | "This Agent can chat now. Bind a Project for repository Work; run `codewhale account projects list` to find one." |
| 1852 | )?; |
| 1853 | } |
| 1854 | Ok(()) |
| 1855 | } |
| 1856 | CloudAgentsCommand::BindProject { agent, project_id } => { |
| 1857 | let listing: AgentListResponse = serde_json::from_value(client.agents()?) |
| 1858 | .context("The Codewhale service returned an invalid Agent list")?; |
| 1859 | let selected = resolve_account_agent(&listing.agents, &agent)?; |
| 1860 | let project_id = validate_resource_id(&project_id, "Project")?; |
| 1861 | let projects: ProjectListResponse = serde_json::from_value(client.projects()?) |
| 1862 | .context("The Codewhale service returned an invalid Project list")?; |
| 1863 | if !projects |
| 1864 | .projects |
| 1865 | .iter() |
| 1866 | .any(|project| project.id == project_id) |
| 1867 | { |
| 1868 | bail!( |
| 1869 | "Project {project_id} is not available on this account. Run `codewhale account projects list`" |
| 1870 | ); |
| 1871 | } |
| 1872 | if selected.project_id == project_id { |
| 1873 | writeln!( |
| 1874 | out, |
| 1875 | "Agent {} is already bound to Project {project_id}.", |
| 1876 | printable(&selected.name) |
| 1877 | )?; |
| 1878 | return Ok(()); |
| 1879 | } |
| 1880 | let bound = client.bind_agent_project(&selected.id, project_id, selected.revision)?; |
| 1881 | writeln!( |
| 1882 | out, |
| 1883 | "Agent {} is bound to Project {project_id}.", |
| 1884 | printable(&bound.name) |
| 1885 | )?; |
| 1886 | Ok(()) |
| 1887 | } |
| 1888 | CloudAgentsCommand::Threads { agent, json } => { |
| 1889 | let listing: AgentListResponse = serde_json::from_value(client.agents()?) |
| 1890 | .context("The Codewhale service returned an invalid Agent list")?; |
| 1891 | let selected = resolve_account_agent(&listing.agents, &agent)?; |
| 1892 | let threads = client |
| 1893 | .threads(&selected.id)? |
| 1894 | .into_iter() |
| 1895 | .filter(|thread| active_agent_thread(thread, &selected.id)) |
| 1896 | .collect::<Vec<_>>(); |
| 1897 | if json { |
| 1898 | return write_computer_json(out, &serde_json::to_value(&threads)?); |
| 1899 | } |
| 1900 | writeln!( |
| 1901 | out, |
| 1902 | "Recent conversations for {} ({})", |
| 1903 | printable(&selected.name), |
| 1904 | threads.len() |
| 1905 | )?; |
| 1906 | for thread in &threads { |
| 1907 | write_agent_thread(out, thread)?; |
| 1908 | } |
| 1909 | Ok(()) |
| 1910 | } |
| 1911 | CloudAgentsCommand::NewThread { |
| 1912 | agent, |
| 1913 | title, |
| 1914 | provider, |
| 1915 | model, |
| 1916 | operation_key, |
| 1917 | } => { |
| 1918 | let listing: AgentListResponse = serde_json::from_value(client.agents()?) |
| 1919 | .context("The Codewhale service returned an invalid Agent list")?; |
| 1920 | let selected = resolve_account_agent(&listing.agents, &agent)?; |
| 1921 | let project_id = selected.project_id.trim(); |
| 1922 | // The saved route is a promise about which model answers, so a |
| 1923 | // route the live catalog does not list is refused before anything |
| 1924 | // is created rather than accepted and discovered at first send. |
| 1925 | let (route_provider, route_model) = validate_model_route(&provider, &model)?; |
| 1926 | let catalog = client.provider_catalog()?; |
| 1927 | let route_provider = assert_catalog_route(&catalog, route_provider, route_model)?; |
| 1928 | let thread = client.create_agent_thread( |
| 1929 | &selected.id, |
| 1930 | if project_id.is_empty() { |
| 1931 | None |
| 1932 | } else { |
| 1933 | Some(project_id) |
| 1934 | }, |
| 1935 | &title, |
| 1936 | route_provider, |
| 1937 | route_model, |
| 1938 | &operation_key, |
| 1939 | )?; |
| 1940 | if !active_agent_thread(&thread, &selected.id) { |
| 1941 | bail!("The Codewhale service did not return an active conversation for this Agent"); |
| 1942 | } |
| 1943 | if thread.model_provider != route_provider || thread.model != route_model { |
| 1944 | bail!( |
| 1945 | "The Codewhale service created conversation {} on {}/{} instead of the requested {route_provider}/{route_model}. Do not send to it; create a new conversation or report this route mismatch", |
| 1946 | printable(&thread.id), |
| 1947 | printable(&thread.model_provider), |
| 1948 | printable(&thread.model) |
| 1949 | ); |
| 1950 | } |
| 1951 | if !project_id.is_empty() && thread.project_id != project_id { |
| 1952 | bail!( |
| 1953 | "The Codewhale service created this conversation in a different Project; verify the Agent's Project before sending" |
| 1954 | ); |
| 1955 | } |
| 1956 | write_agent_thread(out, &thread)?; |
| 1957 | writeln!( |
| 1958 | out, |
| 1959 | "Create request ID: {}", |
| 1960 | validate_operation_key(&operation_key)? |
| 1961 | )?; |
| 1962 | writeln!( |
| 1963 | out, |
| 1964 | "Sending requires an explicit --billing-mode and a new --operation-key for each message." |
| 1965 | )?; |
| 1966 | Ok(()) |
| 1967 | } |
| 1968 | CloudAgentsCommand::Send { |
| 1969 | agent, |
| 1970 | prompt, |
| 1971 | thread, |
| 1972 | billing_mode, |
| 1973 | operation_key, |
| 1974 | } => { |
| 1975 | let listing: AgentListResponse = serde_json::from_value(client.agents()?) |
| 1976 | .context("The Codewhale service returned an invalid Agent list")?; |
| 1977 | let selected = resolve_account_agent(&listing.agents, &agent)?; |
| 1978 | let selected_thread = if let Some(id) = thread { |
| 1979 | client.thread(&id)? |
| 1980 | } else { |
| 1981 | let threads = client.threads(&selected.id)?; |
| 1982 | if threads.len() >= 100 { |
| 1983 | bail!( |
| 1984 | "Agent {} has at least 100 recent conversations; pass --thread with an ID so an older Main conversation is not missed", |
| 1985 | printable(&selected.name) |
| 1986 | ); |
| 1987 | } |
| 1988 | let mut active = threads |
| 1989 | .into_iter() |
| 1990 | .filter(|thread| active_agent_thread(thread, &selected.id)) |
| 1991 | .collect::<Vec<_>>(); |
| 1992 | if active.len() == 1 { |
| 1993 | active.remove(0) |
| 1994 | } else { |
| 1995 | let mut main = active.into_iter().filter(|thread| thread.title == "Main"); |
| 1996 | let first = main.next().ok_or_else(|| anyhow!( |
| 1997 | "Agent {} has no unique active conversation. Run `codewhale account agents new-thread` first, or pass --thread", |
| 1998 | printable(&selected.name) |
| 1999 | ))?; |
| 2000 | if main.next().is_some() { |
| 2001 | bail!( |
| 2002 | "Agent {} has several Main conversations; pass --thread with an ID", |
| 2003 | printable(&selected.name) |
| 2004 | ); |
| 2005 | } |
| 2006 | first |
| 2007 | } |
| 2008 | }; |
| 2009 | if !active_agent_thread(&selected_thread, &selected.id) { |
| 2010 | bail!( |
| 2011 | "That conversation is not active or does not belong to Agent {}", |
| 2012 | printable(&selected.name) |
| 2013 | ); |
| 2014 | } |
| 2015 | if !selected.project_id.is_empty() && selected_thread.project_id != selected.project_id |
| 2016 | { |
| 2017 | bail!( |
| 2018 | "That conversation belongs to a different Project than Agent {} currently owns. Bind the Agent and create a new conversation in its Project before sending", |
| 2019 | printable(&selected.name) |
| 2020 | ); |
| 2021 | } |
| 2022 | let receipt = |
| 2023 | client.send_agent_turn(&selected_thread, &prompt, &billing_mode, &operation_key)?; |
| 2024 | validate_turn_id(&receipt.id)?; |
| 2025 | writeln!(out, "Conversation ID: {}", selected_thread.id)?; |
| 2026 | writeln!(out, "Turn ID: {}", receipt.id)?; |
| 2027 | writeln!(out, "Status: {}", printable(&receipt.status))?; |
| 2028 | writeln!( |
| 2029 | out, |
| 2030 | "Message request ID: {}", |
| 2031 | validate_operation_key(&operation_key)? |
| 2032 | )?; |
| 2033 | writeln!( |
| 2034 | out, |
| 2035 | "Read the answer: codewhale account agents result {} {} {}", |
| 2036 | selected.id, selected_thread.id, receipt.id |
| 2037 | )?; |
| 2038 | Ok(()) |
| 2039 | } |
| 2040 | CloudAgentsCommand::Result { |
| 2041 | agent, |
| 2042 | thread, |
| 2043 | turn, |
| 2044 | since_seq, |
| 2045 | } => { |
| 2046 | let listing: AgentListResponse = serde_json::from_value(client.agents()?) |
| 2047 | .context("The Codewhale service returned an invalid Agent list")?; |
| 2048 | let selected = resolve_account_agent(&listing.agents, &agent)?; |
| 2049 | let selected_thread = client.thread(&thread)?; |
| 2050 | if !active_agent_thread(&selected_thread, &selected.id) { |
| 2051 | bail!( |
| 2052 | "That conversation is not active or does not belong to Agent {}", |
| 2053 | printable(&selected.name) |
| 2054 | ); |
| 2055 | } |
| 2056 | let turn_id = validate_turn_id(&turn)?; |
| 2057 | let result = client.agent_turn_result(&selected_thread.id, turn_id, since_seq)?; |
| 2058 | writeln!(out, "Turn ID: {turn_id}")?; |
| 2059 | writeln!(out, "Status: {}", result.status)?; |
| 2060 | writeln!(out, "Last event sequence: {}", result.last_seq)?; |
| 2061 | if !result.seen_turn { |
| 2062 | writeln!(out, "No events for this turn after the selected sequence.")?; |
| 2063 | } else if !result.answer.is_empty() { |
| 2064 | let label = if result.status == "pending" { |
| 2065 | "Answer so far" |
| 2066 | } else { |
| 2067 | "Answer" |
| 2068 | }; |
| 2069 | writeln!(out, "{label}:")?; |
| 2070 | writeln!( |
| 2071 | out, |
| 2072 | "{}", |
| 2073 | result |
| 2074 | .answer |
| 2075 | .chars() |
| 2076 | .filter(|ch| *ch == '\n' || *ch == '\t' || !ch.is_control()) |
| 2077 | .collect::<String>() |
| 2078 | )?; |
| 2079 | } |
| 2080 | if result.status == "pending" || result.status == "not_seen" { |
| 2081 | writeln!(out, "Run this result command again to read later events.")?; |
| 2082 | } |
| 2083 | Ok(()) |
| 2084 | } |
| 2085 | CloudAgentsCommand::Work { |
| 2086 | agent, |
| 2087 | objective, |
| 2088 | message_id, |
| 2089 | } => { |
| 2090 | let listing: AgentListResponse = serde_json::from_value(client.agents()?) |
| 2091 | .context("The Codewhale service returned an invalid Agent list")?; |
| 2092 | let selected = resolve_account_agent(&listing.agents, &agent)?; |
| 2093 | if selected.project_id.is_empty() { |
| 2094 | bail!( |
| 2095 | "Agent {} needs a repository Project before it can do Work", |
| 2096 | printable(&selected.name) |
| 2097 | ); |
| 2098 | } |
| 2099 | work::assign(client, out, &selected, &objective, &message_id) |
| 2100 | } |
| 2101 | CloudAgentsCommand::WorkStatus { id } => { |
| 2102 | let result = client.agent_work_status(&id)?; |
| 2103 | let run = result |
| 2104 | .get("run") |
| 2105 | .filter(|run| run.is_object()) |
| 2106 | .ok_or_else(|| anyhow!("The Codewhale service returned no Work record"))?; |
| 2107 | if run.get("id").and_then(|value| value.as_str()) |
| 2108 | != Some(validate_resource_id(&id, "Work")?) |
| 2109 | { |
| 2110 | bail!("The Codewhale service returned a different Work record"); |
| 2111 | } |
| 2112 | writeln!(out, "Work ID: {}", validate_resource_id(&id, "Work")?)?; |
| 2113 | writeln!( |
| 2114 | out, |
| 2115 | "Status: {}", |
| 2116 | printable( |
| 2117 | run.get("state") |
| 2118 | .and_then(|value| value.as_str()) |
| 2119 | .unwrap_or("unknown") |
| 2120 | ) |
| 2121 | )?; |
| 2122 | if let Some(title) = run.get("title").and_then(|value| value.as_str()) { |
| 2123 | writeln!(out, "Objective: {}", printable(title))?; |
| 2124 | } |
| 2125 | Ok(()) |
| 2126 | } |
| 2127 | CloudAgentsCommand::WorkCancel { id, queue, reason } => { |
| 2128 | work::cancel(client, out, &id, queue, reason.as_deref()) |
| 2129 | } |
| 2130 | CloudAgentsCommand::WorkResult { id, json } => work::result(client, out, &id, json), |
| 2131 | CloudAgentsCommand::WorkQuote { id, operation_key } => { |
| 2132 | work::quote(client, out, &id, &operation_key) |
| 2133 | } |
| 2134 | CloudAgentsCommand::WorkLaunch { |
| 2135 | id, |
| 2136 | operation_key, |
| 2137 | confirmation, |
| 2138 | confirm_eu_compute, |
| 2139 | } => { |
| 2140 | let confirmation = if confirmation == "-" && confirm_eu_compute { |
| 2141 | if io::stdin().is_terminal() { |
| 2142 | bail!( |
| 2143 | "Pipe the launch confirmation to stdin; it must not be typed into the terminal" |
| 2144 | ); |
| 2145 | } |
| 2146 | work::read_confirmation(io::stdin().lock())? |
| 2147 | } else { |
| 2148 | confirmation |
| 2149 | }; |
| 2150 | work::launch( |
| 2151 | client, |
| 2152 | out, |
| 2153 | &id, |
| 2154 | &operation_key, |
| 2155 | &confirmation, |
| 2156 | confirm_eu_compute, |
| 2157 | ) |
| 2158 | } |
| 2159 | } |
| 2160 | } |
| 2161 | |
| 2162 | fn write_computer<W: Write>(out: &mut W, computer: &AccountComputer) -> Result<()> { |
| 2163 | validate_computer_id(&computer.id)?; |
| 2164 | if computer.owner_id.trim().is_empty() { |
| 2165 | bail!("The Codewhale service returned a Computer without an owner"); |
| 2166 | } |
| 2167 | writeln!(out, "Computer: {}", printable(&computer.name))?; |
| 2168 | writeln!(out, "ID: {}", computer.id)?; |
| 2169 | writeln!(out, "Account ID: {}", printable(&computer.owner_id))?; |
| 2170 | writeln!(out, "Status: {}", printable(&computer.status))?; |
| 2171 | writeln!(out, "Region: {}", printable(&computer.region))?; |
| 2172 | Ok(()) |
| 2173 | } |
| 2174 | |
| 2175 | fn write_computer_json<W: Write>(out: &mut W, value: &serde_json::Value) -> Result<()> { |
| 2176 | serde_json::to_writer_pretty(&mut *out, value) |
| 2177 | .context("failed to write Codewhale Computer JSON")?; |
| 2178 | writeln!(out)?; |
| 2179 | Ok(()) |
| 2180 | } |
| 2181 | |
| 2182 | fn run_computers<T: CloudTransport, W: Write>( |
| 2183 | command: CloudComputersCommand, |
| 2184 | client: &CloudClient<'_, T>, |
| 2185 | machine: &machine::MachineKeyEnv, |
| 2186 | out: &mut W, |
| 2187 | ) -> Result<()> { |
| 2188 | // Machine keys are intentionally narrower than an interactive account |
| 2189 | // session. A present key must never silently fall back to a human login. |
| 2190 | if machine.is_present() { |
| 2191 | bail!( |
| 2192 | "Computers require an interactive Codewhale account login; unset CODEWHALE_API_KEY and run `codewhale login`" |
| 2193 | ); |
| 2194 | } |
| 2195 | match command { |
| 2196 | CloudComputersCommand::List { json } => { |
| 2197 | let response = client.computers()?; |
| 2198 | if json { |
| 2199 | return write_computer_json(out, &response); |
| 2200 | } |
| 2201 | let computers: ComputerListResponse = serde_json::from_value(response) |
| 2202 | .context("The Codewhale service returned an invalid Computer list")?; |
| 2203 | let computers = computers.computers; |
| 2204 | writeln!(out, "Codewhale Computers ({})", computers.len())?; |
| 2205 | for computer in &computers { |
| 2206 | write_computer(out, computer)?; |
| 2207 | } |
| 2208 | Ok(()) |
| 2209 | } |
| 2210 | CloudComputersCommand::Create { |
| 2211 | name, |
| 2212 | boat_trial, |
| 2213 | eu_compute_opt_in, |
| 2214 | } => { |
| 2215 | let computer = client.create_computer(&name, boat_trial, eu_compute_opt_in)?; |
| 2216 | write_computer(out, &computer)?; |
| 2217 | writeln!( |
| 2218 | out, |
| 2219 | "Saved Computer identity; compute is allocated when you start it." |
| 2220 | )?; |
| 2221 | Ok(()) |
| 2222 | } |
| 2223 | CloudComputersCommand::Show { id, json } => { |
| 2224 | let response = client.computer(&id)?; |
| 2225 | if json { |
| 2226 | write_computer_json(out, &response) |
| 2227 | } else { |
| 2228 | let result: ComputerResponse = serde_json::from_value(response) |
| 2229 | .context("The Codewhale service returned an invalid Computer")?; |
| 2230 | write_computer(out, &result.computer) |
| 2231 | } |
| 2232 | } |
| 2233 | CloudComputersCommand::Usage { id, json } => { |
| 2234 | if !json { |
| 2235 | bail!("Use --json to inspect Computer usage receipts"); |
| 2236 | } |
| 2237 | write_computer_json(out, &client.computer_usage(&id)?) |
| 2238 | } |
| 2239 | CloudComputersCommand::Start { id } => { |
| 2240 | let result = client.computer_action(&id, "start")?; |
| 2241 | if result.queued { |
| 2242 | writeln!(out, "Computer start queued.")?; |
| 2243 | if !result.computer.start_queue_reason.is_empty() { |
| 2244 | writeln!( |
| 2245 | out, |
| 2246 | "Reason: {}", |
| 2247 | printable(&result.computer.start_queue_reason) |
| 2248 | )?; |
| 2249 | } |
| 2250 | } |
| 2251 | write_computer(out, &result.computer) |
| 2252 | } |
| 2253 | CloudComputersCommand::Pause { id } => { |
| 2254 | let result = client.computer_action(&id, "pause")?; |
| 2255 | write_computer(out, &result.computer) |
| 2256 | } |
| 2257 | CloudComputersCommand::Delete { id } => { |
| 2258 | client.delete_computer(&id)?; |
| 2259 | writeln!(out, "Deleted Computer {}.", validate_computer_id(&id)?)?; |
| 2260 | Ok(()) |
| 2261 | } |
| 2262 | } |
| 2263 | } |
| 2264 | |
| 2265 | enum KeyReadMode { |
| 2266 | Stdin, |
| 2267 | HiddenPrompt(String), |
| 2268 | } |
| 2269 | |
| 2270 | pub(crate) fn run(args: CloudArgs, profile: Option<&str>, config: &ConfigStore) -> Result<()> { |
| 2271 | let machine = machine::MachineKeyEnv::from_process_env(); |
| 2272 | let requested_base = machine::resolve_api_base( |
| 2273 | args.api_base.as_deref(), |
| 2274 | std::env::var(machine::MACHINE_API_BASE_ENV).ok().as_deref(), |
| 2275 | std::env::var(CLOUD_API_BASE_ENV).ok().as_deref(), |
| 2276 | DEFAULT_API_BASE, |
| 2277 | ); |
| 2278 | if machine.is_present() { |
| 2279 | // A machine token is a bearer credential with no replay protection. |
| 2280 | // Refuse cleartext to a remote host before a transport exists, so |
| 2281 | // there is no code path on which the key could be written to a socket. |
| 2282 | machine::require_secure_base(&requested_base)?; |
| 2283 | } |
| 2284 | let api_base = validate_api_base(&requested_base)?; |
| 2285 | let transport = ReqwestTransport::new(api_base.url.clone())?; |
| 2286 | // Account refresh tokens require an OS credential manager. The ordinary |
| 2287 | // provider backend remains independently configurable for `--from-local`. |
| 2288 | let cloud_secrets = cloud_session_secrets()?; |
| 2289 | let provider_secrets = Secrets::auto_detect(); |
| 2290 | let profile = normalized_profile(profile); |
| 2291 | let mut stdout = io::stdout().lock(); |
| 2292 | let mut key_reader = |mode: KeyReadMode| match mode { |
| 2293 | KeyReadMode::Stdin => read_key_from_stdin(), |
| 2294 | KeyReadMode::HiddenPrompt(provider) => read_key_hidden(&provider), |
| 2295 | }; |
| 2296 | let mut opener = |url: String| webbrowser::open(&url).is_ok(); |
| 2297 | let mut sleeper = |duration| thread::sleep(duration); |
| 2298 | run_with( |
| 2299 | args.command, |
| 2300 | &profile, |
| 2301 | &api_base.display, |
| 2302 | config, |
| 2303 | &cloud_secrets, |
| 2304 | &provider_secrets, |
| 2305 | &machine, |
| 2306 | &transport, |
| 2307 | &mut stdout, |
| 2308 | &mut key_reader, |
| 2309 | &mut opener, |
| 2310 | &mut sleeper, |
| 2311 | ) |
| 2312 | } |
| 2313 | |
| 2314 | fn cloud_session_secrets() -> Result<Secrets> { |
| 2315 | // Codex-style storage contract: the OS credential manager is preferred |
| 2316 | // but never required; without one, sessions live in the private 0600 |
| 2317 | // Codewhale secrets file. Only an unresolvable store path fails here. |
| 2318 | secure_account_session_secrets().map_err(|err| anyhow!(err.to_string())) |
| 2319 | } |
| 2320 | |
| 2321 | /// `codewhale login` is a convenience entry to the account device flow — the |
| 2322 | /// same path as `codewhale account login`, without re-spelling the subcommand. |
| 2323 | pub(crate) fn run_account_login( |
| 2324 | no_open: bool, |
| 2325 | timeout_seconds: u64, |
| 2326 | profile: Option<&str>, |
| 2327 | config: &ConfigStore, |
| 2328 | ) -> Result<()> { |
| 2329 | run( |
| 2330 | CloudArgs { |
| 2331 | api_base: None, |
| 2332 | command: CloudCommand::Login(CloudLoginArgs { |
| 2333 | no_open, |
| 2334 | timeout_seconds, |
| 2335 | }), |
| 2336 | }, |
| 2337 | profile, |
| 2338 | config, |
| 2339 | ) |
| 2340 | } |
| 2341 | |
| 2342 | pub(crate) fn reject_inline_api_key(api_key: Option<&str>) -> Result<()> { |
| 2343 | if api_key.is_some() { |
| 2344 | bail!( |
| 2345 | "`codewhale account` does not accept the global `--api-key` flag because command-line values can leak through shell history. Use `account keys set <provider>` for a hidden prompt, `--api-key-stdin`, or `--from-local`" |
| 2346 | ); |
| 2347 | } |
| 2348 | Ok(()) |
| 2349 | } |
| 2350 | |
| 2351 | #[allow(clippy::too_many_arguments)] |
| 2352 | fn run_with<T: CloudTransport, W: Write>( |
| 2353 | command: CloudCommand, |
| 2354 | profile: &str, |
| 2355 | api_base: &str, |
| 2356 | config: &ConfigStore, |
| 2357 | cloud_secrets: &Secrets, |
| 2358 | provider_secrets: &Secrets, |
| 2359 | machine: &machine::MachineKeyEnv, |
| 2360 | transport: &T, |
| 2361 | out: &mut W, |
| 2362 | key_reader: &mut dyn FnMut(KeyReadMode) -> Result<String>, |
| 2363 | opener: &mut dyn FnMut(String) -> bool, |
| 2364 | sleeper: &mut dyn FnMut(Duration), |
| 2365 | ) -> Result<()> { |
| 2366 | let client = CloudClient::new(transport, cloud_secrets, profile, api_base); |
| 2367 | match command { |
| 2368 | CloudCommand::Login(login) => { |
| 2369 | let device = client.start_device()?; |
| 2370 | validate_user_code(&device.user_code)?; |
| 2371 | validate_verification_url( |
| 2372 | &device.verification_uri, |
| 2373 | api_base, |
| 2374 | &device.user_code, |
| 2375 | false, |
| 2376 | )?; |
| 2377 | let verification_uri_complete = validate_verification_url( |
| 2378 | &device.verification_uri_complete, |
| 2379 | api_base, |
| 2380 | &device.user_code, |
| 2381 | true, |
| 2382 | )?; |
| 2383 | writeln!(out, "Codewhale account sign-in")?; |
| 2384 | writeln!(out, "Code: {}", device.user_code)?; |
| 2385 | writeln!(out, "Open: {verification_uri_complete}")?; |
| 2386 | writeln!(out, "Profile: {}", printable(profile))?; |
| 2387 | if !login.no_open && !opener(verification_uri_complete) { |
| 2388 | writeln!( |
| 2389 | out, |
| 2390 | "Browser could not be opened; use the URL and code above." |
| 2391 | )?; |
| 2392 | } |
| 2393 | let _ = |
| 2394 | client.poll_device(&device, Duration::from_secs(login.timeout_seconds), sleeper)?; |
| 2395 | let user = client.me()?; |
| 2396 | write_account(out, "Signed in to Codewhale.", profile, api_base, &user)?; |
| 2397 | Ok(()) |
| 2398 | } |
| 2399 | CloudCommand::Status => match client.load_auth()? { |
| 2400 | Some(_) => { |
| 2401 | let user = client.me()?; |
| 2402 | write_account(out, "Signed in to Codewhale.", profile, api_base, &user) |
| 2403 | } |
| 2404 | None => { |
| 2405 | writeln!(out, "Not signed in to Codewhale.")?; |
| 2406 | writeln!(out, "Profile: {}", printable(profile))?; |
| 2407 | writeln!(out, "API: {api_base}")?; |
| 2408 | writeln!(out, "Run `codewhale login` to sign in.")?; |
| 2409 | Ok(()) |
| 2410 | } |
| 2411 | }, |
| 2412 | CloudCommand::Logout => { |
| 2413 | let remote_revoked = client.logout()?; |
| 2414 | writeln!(out, "Removed the local Codewhale account session.")?; |
| 2415 | writeln!(out, "Profile: {}", printable(profile))?; |
| 2416 | if !remote_revoked { |
| 2417 | writeln!( |
| 2418 | out, |
| 2419 | "Remote revocation was not confirmed; the local tokens are gone." |
| 2420 | )?; |
| 2421 | } |
| 2422 | Ok(()) |
| 2423 | } |
| 2424 | CloudCommand::Keys(keys) => match keys.command { |
| 2425 | CloudKeysCommand::List => { |
| 2426 | let user = client.me()?; |
| 2427 | let catalog = client.provider_catalog()?; |
| 2428 | write_account(out, "Codewhale account keys.", profile, api_base, &user)?; |
| 2429 | for row in &catalog { |
| 2430 | let stored = user.model_keys.get(&row.id); |
| 2431 | let status = match stored { |
| 2432 | Some(state) if state.configured => match state |
| 2433 | .state |
| 2434 | .as_deref() |
| 2435 | .map(printable) |
| 2436 | .filter(|value| !value.is_empty()) |
| 2437 | { |
| 2438 | Some(reported) => format!("set ({reported})"), |
| 2439 | None => "set".to_string(), |
| 2440 | }, |
| 2441 | _ => "not set".to_string(), |
| 2442 | }; |
| 2443 | writeln!(out, "{}: {status} — {}", row.id, row.display_label())?; |
| 2444 | } |
| 2445 | Ok(()) |
| 2446 | } |
| 2447 | CloudKeysCommand::Set(set) => { |
| 2448 | let provider = validate_provider_id(&set.provider)?; |
| 2449 | let user = client.me()?; |
| 2450 | let catalog = client.provider_catalog()?; |
| 2451 | let row = catalog_row(&catalog, &provider)?; |
| 2452 | let key = if set.from_local { |
| 2453 | let kind = row.local_kind().ok_or_else(|| { |
| 2454 | anyhow!( |
| 2455 | "`{provider}` has no local runtime provider, so there is no local key to copy. Use `--api-key-stdin` or the hidden prompt" |
| 2456 | ) |
| 2457 | })?; |
| 2458 | resolve_local_key(config, provider_secrets, kind)?.ok_or_else(|| { |
| 2459 | anyhow!( |
| 2460 | "No local {} API key was found in config, the secret store, or the environment", |
| 2461 | kind.as_str() |
| 2462 | ) |
| 2463 | })? |
| 2464 | } else if set.api_key_stdin { |
| 2465 | key_reader(KeyReadMode::Stdin)? |
| 2466 | } else { |
| 2467 | key_reader(KeyReadMode::HiddenPrompt(provider.clone()))? |
| 2468 | }; |
| 2469 | let key = key.trim().to_string(); |
| 2470 | validate_api_key(&key)?; |
| 2471 | let label = validate_label(&set.label)?; |
| 2472 | client.set_key(&provider, &key, &label)?; |
| 2473 | writeln!( |
| 2474 | out, |
| 2475 | "Saved {provider} for Codewhale account {} (profile {}).", |
| 2476 | printable(&user.id), |
| 2477 | printable(profile) |
| 2478 | )?; |
| 2479 | Ok(()) |
| 2480 | } |
| 2481 | CloudKeysCommand::Remove { provider } => { |
| 2482 | let provider = validate_provider_id(&provider)?; |
| 2483 | let user = client.me()?; |
| 2484 | let catalog = client.provider_catalog()?; |
| 2485 | let _ = catalog_row(&catalog, &provider)?; |
| 2486 | client.remove_key(&provider)?; |
| 2487 | writeln!( |
| 2488 | out, |
| 2489 | "Removed {provider} from Codewhale account {} (profile {}).", |
| 2490 | printable(&user.id), |
| 2491 | printable(profile) |
| 2492 | )?; |
| 2493 | Ok(()) |
| 2494 | } |
| 2495 | }, |
| 2496 | CloudCommand::Computers(computers) => { |
| 2497 | run_computers(computers.command, &client, machine, out) |
| 2498 | } |
| 2499 | CloudCommand::Projects(projects) => run_projects(projects.command, &client, machine, out), |
| 2500 | CloudCommand::Github(github) => run_github(github.command, &client, machine, out), |
| 2501 | CloudCommand::Agents(agents) => run_agents(agents.command, &client, machine, out), |
| 2502 | CloudCommand::ApiKeys(api_keys) => { |
| 2503 | machine::run_api_keys(api_keys, &client, machine, provider_secrets, out, sleeper) |
| 2504 | } |
| 2505 | CloudCommand::Whoami => match machine.resolve()? { |
| 2506 | // A present machine key wins and never falls back: silently |
| 2507 | // downgrading a machine credential to a human one is how CI ends |
| 2508 | // up running as the wrong identity. |
| 2509 | Some(key) => { |
| 2510 | let machine_client = machine::MachineClient::new(transport, key); |
| 2511 | let who = machine_client.whoami(sleeper)?; |
| 2512 | machine::write_whoami(out, &who, api_base, machine_client.key_head()) |
| 2513 | } |
| 2514 | None => { |
| 2515 | let user = client.me()?; |
| 2516 | write_account(out, "Signed in to Codewhale.", profile, api_base, &user) |
| 2517 | } |
| 2518 | }, |
| 2519 | CloudCommand::Agent => { |
| 2520 | // The agent route is machine-key-only by design, so CI and humans |
| 2521 | // never blur in an audit trail. There is no session fallback. |
| 2522 | let key = machine.require()?; |
| 2523 | let machine_client = machine::MachineClient::new(transport, key); |
| 2524 | let agent = machine_client.agent(sleeper)?.agent; |
| 2525 | machine::write_agent(out, &agent) |
| 2526 | } |
| 2527 | CloudCommand::Pull(args) => { |
| 2528 | if !args.dry_run { |
| 2529 | bail!( |
| 2530 | "Account settings import is not available yet; local config was not changed. Run `codewhale account pull --dry-run` to inspect the signed-in account." |
| 2531 | ); |
| 2532 | } |
| 2533 | let user = client.me()?; |
| 2534 | // `/api/me` currently exposes account identity and key metadata, |
| 2535 | // not a versioned settings document that can be applied locally. |
| 2536 | // Stay read-only and explicit until that import contract exists. |
| 2537 | writeln!(out, "Account settings (pull --dry-run):")?; |
| 2538 | writeln!(out, "Account ID: {}", printable(&user.id))?; |
| 2539 | writeln!(out, "Profile: {}", printable(profile))?; |
| 2540 | writeln!(out, "API: {api_base}")?; |
| 2541 | writeln!( |
| 2542 | out, |
| 2543 | "dry-run: remote settings import is not available; local config unchanged" |
| 2544 | )?; |
| 2545 | // Show the invariant: Bearer custody stays in the OS keyring, never in config.toml. |
| 2546 | writeln!( |
| 2547 | out, |
| 2548 | "Secure custody: Bearer tokens remain in the OS keyring" |
| 2549 | )?; |
| 2550 | Ok(()) |
| 2551 | } |
| 2552 | CloudCommand::Push(args) => { |
| 2553 | let user = client.me()?; |
| 2554 | if !args.dry_run { |
| 2555 | bail!( |
| 2556 | "Push is never automatic; re-run with --dry-run to preview, then confirm explicitly" |
| 2557 | ); |
| 2558 | } |
| 2559 | writeln!(out, "Account settings (push --dry-run):")?; |
| 2560 | writeln!(out, "Account ID: {}", printable(&user.id))?; |
| 2561 | writeln!(out, "Profile: {}", printable(profile))?; |
| 2562 | writeln!(out, "API: {api_base}")?; |
| 2563 | writeln!( |
| 2564 | out, |
| 2565 | "dry-run: would PATCH /api/me/preferences with If-Match revision check (412 on conflict)" |
| 2566 | )?; |
| 2567 | writeln!( |
| 2568 | out, |
| 2569 | "No credentials, paths, or env are copied; only explicit fields (field-level last-writer-wins)" |
| 2570 | )?; |
| 2571 | Ok(()) |
| 2572 | } |
| 2573 | } |
| 2574 | } |
| 2575 | |
| 2576 | /// Resolve the account's configured agent route for a machine-key run. |
| 2577 | /// |
| 2578 | /// `codewhale review` hard-errors when a model resolves to several configured |
| 2579 | /// routes. When CI authenticates with a machine key, the account has already |
| 2580 | /// answered that question, so its configured provider is the disambiguator — |
| 2581 | /// no new flag, and no guess. Returns `None` when no machine key is set, which |
| 2582 | /// leaves the ordinary local resolution untouched. |
| 2583 | pub(crate) fn machine_review_provider() -> Result<Option<ProviderKind>> { |
| 2584 | let machine = machine::MachineKeyEnv::from_process_env(); |
| 2585 | let Some(key) = machine.resolve()? else { |
| 2586 | return Ok(None); |
| 2587 | }; |
| 2588 | let requested_base = machine::resolve_api_base( |
| 2589 | None, |
| 2590 | std::env::var(machine::MACHINE_API_BASE_ENV).ok().as_deref(), |
| 2591 | std::env::var(CLOUD_API_BASE_ENV).ok().as_deref(), |
| 2592 | DEFAULT_API_BASE, |
| 2593 | ); |
| 2594 | machine::require_secure_base(&requested_base)?; |
| 2595 | let api_base = validate_api_base(&requested_base)?; |
| 2596 | let transport = ReqwestTransport::new(api_base.url)?; |
| 2597 | let client = machine::MachineClient::new(&transport, key); |
| 2598 | // The call that actually needs a model is the call that refuses without |
| 2599 | // one, so this precondition runs before any review work starts. |
| 2600 | let agent = client.agent(&mut |duration| thread::sleep(duration))?.agent; |
| 2601 | machine::review_provider_from_agent(&agent).map(Some) |
| 2602 | } |
| 2603 | |
| 2604 | fn write_account<W: Write>( |
| 2605 | out: &mut W, |
| 2606 | heading: &str, |
| 2607 | profile: &str, |
| 2608 | api_base: &str, |
| 2609 | user: &CloudUser, |
| 2610 | ) -> Result<()> { |
| 2611 | writeln!(out, "{heading}")?; |
| 2612 | writeln!(out, "Account ID: {}", printable(&user.id))?; |
| 2613 | if !user.display_name.trim().is_empty() { |
| 2614 | writeln!(out, "Name: {}", printable(&user.display_name))?; |
| 2615 | } |
| 2616 | if !user.email.trim().is_empty() { |
| 2617 | writeln!(out, "Email: {}", printable(&user.email))?; |
| 2618 | } |
| 2619 | if !user.plan.trim().is_empty() { |
| 2620 | writeln!(out, "Plan: {}", printable(&user.plan))?; |
| 2621 | } |
| 2622 | writeln!(out, "Profile: {}", printable(profile))?; |
| 2623 | writeln!(out, "API: {api_base}")?; |
| 2624 | Ok(()) |
| 2625 | } |
| 2626 | |
| 2627 | struct ValidatedApiBase { |
| 2628 | url: Url, |
| 2629 | display: String, |
| 2630 | } |
| 2631 | |
| 2632 | fn validate_api_base(value: &str) -> Result<ValidatedApiBase> { |
| 2633 | let mut url = Url::parse(value.trim()).context("invalid Codewhale account API base URL")?; |
| 2634 | if !url.username().is_empty() || url.password().is_some() { |
| 2635 | bail!("Codewhale account API base URL must not contain credentials"); |
| 2636 | } |
| 2637 | if url.query().is_some() || url.fragment().is_some() { |
| 2638 | bail!("Codewhale account API base URL must not contain a query or fragment"); |
| 2639 | } |
| 2640 | if !matches!(url.path(), "" | "/") { |
| 2641 | bail!("Codewhale account API base URL must be an origin without a path"); |
| 2642 | } |
| 2643 | let host = url |
| 2644 | .host_str() |
| 2645 | .ok_or_else(|| anyhow!("Codewhale account API base URL must include a host"))?; |
| 2646 | let allowed = url.scheme() == "https" || (url.scheme() == "http" && is_loopback_host(host)); |
| 2647 | if !allowed { |
| 2648 | bail!( |
| 2649 | "Codewhale account API base URL must use HTTPS (loopback HTTP is allowed for testing)" |
| 2650 | ); |
| 2651 | } |
| 2652 | url.set_path("/"); |
| 2653 | let display = url.as_str().trim_end_matches('/').to_string(); |
| 2654 | Ok(ValidatedApiBase { url, display }) |
| 2655 | } |
| 2656 | |
| 2657 | fn validate_verification_url( |
| 2658 | value: &str, |
| 2659 | api_base: &str, |
| 2660 | user_code: &str, |
| 2661 | complete: bool, |
| 2662 | ) -> Result<String> { |
| 2663 | let url = |
| 2664 | Url::parse(value).context("The Codewhale service returned an invalid verification URL")?; |
| 2665 | if value != url.as_str() { |
| 2666 | bail!("The Codewhale service returned an unsafe verification URL"); |
| 2667 | } |
| 2668 | let host = url.host_str().ok_or_else(|| { |
| 2669 | anyhow!("The Codewhale service returned a verification URL without a host") |
| 2670 | })?; |
| 2671 | if !url.username().is_empty() || url.password().is_some() || url.fragment().is_some() { |
| 2672 | bail!("The Codewhale service returned an unsafe verification URL"); |
| 2673 | } |
| 2674 | if url.path() != "/cli/authorize" { |
| 2675 | bail!("The Codewhale service returned an unsafe verification URL"); |
| 2676 | } |
| 2677 | |
| 2678 | let api = Url::parse(api_base).context("invalid Codewhale account API base URL")?; |
| 2679 | let canonical_api = api.scheme() == "https" |
| 2680 | && api.host_str() == Some("api.codewhale.net") |
| 2681 | && api.port_or_known_default() == Some(443); |
| 2682 | let loopback_api = api.host_str().is_some_and(is_loopback_host); |
| 2683 | if canonical_api { |
| 2684 | if url.scheme() != "https" |
| 2685 | || !host.eq_ignore_ascii_case("app.codewhale.net") |
| 2686 | || url.port_or_known_default() != Some(443) |
| 2687 | { |
| 2688 | bail!("The Codewhale service returned an untrusted verification origin"); |
| 2689 | } |
| 2690 | } else if loopback_api { |
| 2691 | if !matches!(url.scheme(), "http" | "https") || !is_loopback_host(host) { |
| 2692 | bail!("The Codewhale service returned an untrusted verification origin"); |
| 2693 | } |
| 2694 | } else { |
| 2695 | bail!( |
| 2696 | "Browser login is only enabled for the canonical Codewhale account API or a loopback test API" |
| 2697 | ); |
| 2698 | } |
| 2699 | |
| 2700 | let query = url.query_pairs().collect::<Vec<_>>(); |
| 2701 | if complete { |
| 2702 | if query.len() != 1 || query[0].0 != "user_code" || query[0].1 != user_code { |
| 2703 | bail!("The Codewhale service returned an unsafe verification URL"); |
| 2704 | } |
| 2705 | } else if !query.is_empty() { |
| 2706 | bail!("The Codewhale service returned an unsafe verification URL"); |
| 2707 | } |
| 2708 | Ok(url.to_string()) |
| 2709 | } |
| 2710 | |
| 2711 | fn is_loopback_host(host: &str) -> bool { |
| 2712 | let host = host |
| 2713 | .strip_prefix('[') |
| 2714 | .and_then(|value| value.strip_suffix(']')) |
| 2715 | .unwrap_or(host); |
| 2716 | host.eq_ignore_ascii_case("localhost") |
| 2717 | || host |
| 2718 | .parse::<IpAddr>() |
| 2719 | .is_ok_and(|address| address.is_loopback()) |
| 2720 | } |
| 2721 | |
| 2722 | fn validate_user_code(code: &str) -> Result<()> { |
| 2723 | const ALPHABET: &[u8] = b"ABCDEFGHJKLMNPQRSTUVWXYZ23456789"; |
| 2724 | let bytes = code.as_bytes(); |
| 2725 | if bytes.len() != 14 |
| 2726 | || bytes[4] != b'-' |
| 2727 | || bytes[9] != b'-' |
| 2728 | || bytes |
| 2729 | .iter() |
| 2730 | .enumerate() |
| 2731 | .any(|(index, byte)| !matches!(index, 4 | 9) && !ALPHABET.contains(byte)) |
| 2732 | { |
| 2733 | bail!("The Codewhale service returned an invalid user code"); |
| 2734 | } |
| 2735 | Ok(()) |
| 2736 | } |
| 2737 | |
| 2738 | fn validate_device_code(code: &str) -> Result<()> { |
| 2739 | if code.len() != 43 |
| 2740 | || !code |
| 2741 | .bytes() |
| 2742 | .all(|byte| byte.is_ascii_alphanumeric() || matches!(byte, b'-' | b'_')) |
| 2743 | { |
| 2744 | bail!("The Codewhale service returned an invalid device authorization response"); |
| 2745 | } |
| 2746 | Ok(()) |
| 2747 | } |
| 2748 | |
| 2749 | fn validate_api_key(key: &str) -> Result<()> { |
| 2750 | let bytes = key.len(); |
| 2751 | if bytes < MIN_API_KEY_BYTES || bytes as u64 > MAX_API_KEY_BYTES { |
| 2752 | bail!("API key must be {MIN_API_KEY_BYTES}-{MAX_API_KEY_BYTES} UTF-8 bytes"); |
| 2753 | } |
| 2754 | if key.chars().any(is_ascii_control) { |
| 2755 | bail!("API key contains invalid control characters"); |
| 2756 | } |
| 2757 | Ok(()) |
| 2758 | } |
| 2759 | |
| 2760 | fn validate_label(label: &str) -> Result<String> { |
| 2761 | let label = label.split_whitespace().collect::<Vec<_>>().join(" "); |
| 2762 | if label.is_empty() |
| 2763 | || label.chars().count() > MAX_KEY_LABEL_CHARS |
| 2764 | || label.chars().any(is_ascii_control) |
| 2765 | { |
| 2766 | bail!("key label must contain 1-{MAX_KEY_LABEL_CHARS} characters"); |
| 2767 | } |
| 2768 | Ok(label) |
| 2769 | } |
| 2770 | |
| 2771 | fn is_ascii_control(character: char) -> bool { |
| 2772 | character <= '\u{001f}' || character == '\u{007f}' |
| 2773 | } |
| 2774 | |
| 2775 | /// Find one catalog row by id, or fail naming what the account does offer. |
| 2776 | /// |
| 2777 | /// A catalog miss is the common typo, so the message lists the ids rather than |
| 2778 | /// leaving the user to guess or read a 404. |
| 2779 | fn catalog_row<'a>(catalog: &'a [CatalogProvider], provider: &str) -> Result<&'a CatalogProvider> { |
| 2780 | catalog |
| 2781 | .iter() |
| 2782 | .find(|row| row.id == provider) |
| 2783 | .ok_or_else(|| { |
| 2784 | let known = catalog |
| 2785 | .iter() |
| 2786 | .map(|row| row.id.as_str()) |
| 2787 | .collect::<Vec<_>>() |
| 2788 | .join(", "); |
| 2789 | anyhow!("`{provider}` is not a provider this Codewhale account can connect. Known providers: {known}") |
| 2790 | }) |
| 2791 | } |
| 2792 | |
| 2793 | fn resolve_local_key( |
| 2794 | config: &ConfigStore, |
| 2795 | secrets: &Secrets, |
| 2796 | kind: ProviderKind, |
| 2797 | ) -> Result<Option<String>> { |
| 2798 | let provider_config = config.config.providers.for_provider(kind); |
| 2799 | let from_config = provider_config.api_key.clone(); |
| 2800 | if let Some(value) = from_config |
| 2801 | .and_then(resolve_config_key_reference) |
| 2802 | .filter(|value| !value.trim().is_empty()) |
| 2803 | { |
| 2804 | return Ok(Some(value)); |
| 2805 | } |
| 2806 | if let Some(value) = secrets |
| 2807 | .get(kind.as_str()) |
| 2808 | .context("failed to read the local provider secret store")? |
| 2809 | .filter(|value| !value.trim().is_empty()) |
| 2810 | { |
| 2811 | return Ok(Some(value)); |
| 2812 | } |
| 2813 | Ok(kind.provider().env_vars().iter().find_map(|name| { |
| 2814 | std::env::var(name) |
| 2815 | .ok() |
| 2816 | .filter(|value| !value.trim().is_empty()) |
| 2817 | })) |
| 2818 | } |
| 2819 | |
| 2820 | fn resolve_config_key_reference(value: String) -> Option<String> { |
| 2821 | let trimmed = value.trim(); |
| 2822 | let Some(variable) = trimmed.strip_prefix('$') else { |
| 2823 | return Some(value); |
| 2824 | }; |
| 2825 | if variable.is_empty() |
| 2826 | || !variable |
| 2827 | .bytes() |
| 2828 | .all(|byte| byte.is_ascii_alphanumeric() || byte == b'_') |
| 2829 | { |
| 2830 | return None; |
| 2831 | } |
| 2832 | std::env::var(variable) |
| 2833 | .ok() |
| 2834 | .filter(|value| !value.trim().is_empty()) |
| 2835 | } |
| 2836 | |
| 2837 | fn read_key_from_stdin() -> Result<String> { |
| 2838 | let mut bytes = Vec::new(); |
| 2839 | io::stdin() |
| 2840 | .take(MAX_API_KEY_STDIN_BYTES + 1) |
| 2841 | .read_to_end(&mut bytes) |
| 2842 | .context("failed to read API key from stdin")?; |
| 2843 | parse_key_input(bytes) |
| 2844 | } |
| 2845 | |
| 2846 | fn parse_key_input(bytes: Vec<u8>) -> Result<String> { |
| 2847 | if bytes.len() as u64 > MAX_API_KEY_STDIN_BYTES { |
| 2848 | bail!("API key input is unexpectedly large"); |
| 2849 | } |
| 2850 | let value = String::from_utf8(bytes).context("API key from stdin is not valid UTF-8")?; |
| 2851 | let value = value.trim().to_string(); |
| 2852 | validate_api_key(&value)?; |
| 2853 | Ok(value) |
| 2854 | } |
| 2855 | |
| 2856 | fn read_key_hidden(provider: &str) -> Result<String> { |
| 2857 | if !io::stdin().is_terminal() { |
| 2858 | bail!("interactive key entry requires a terminal; use `--api-key-stdin` for piped input"); |
| 2859 | } |
| 2860 | let term = console::Term::stderr(); |
| 2861 | term.write_str(&format!("Enter {provider} API key: ")) |
| 2862 | .context("failed to write API key prompt")?; |
| 2863 | let value = term |
| 2864 | .read_secure_line() |
| 2865 | .context("failed to read API key securely")?; |
| 2866 | term.write_line("").ok(); |
| 2867 | let value = value.trim().to_string(); |
| 2868 | validate_api_key(&value)?; |
| 2869 | Ok(value) |
| 2870 | } |
| 2871 | |
| 2872 | fn json_body(value: &impl Serialize) -> Result<Vec<u8>> { |
| 2873 | serde_json::to_vec(value).context("failed to encode Codewhale account request") |
| 2874 | } |
| 2875 | |
| 2876 | fn expect_json<T: DeserializeOwned>(response: CloudResponse, statuses: &[u16]) -> Result<T> { |
| 2877 | if !statuses.contains(&response.status) { |
| 2878 | return Err(response_error(&response)); |
| 2879 | } |
| 2880 | parse_json_body(&response.body) |
| 2881 | } |
| 2882 | |
| 2883 | fn expect_empty(response: CloudResponse, statuses: &[u16]) -> Result<()> { |
| 2884 | if statuses.contains(&response.status) { |
| 2885 | Ok(()) |
| 2886 | } else { |
| 2887 | Err(response_error(&response)) |
| 2888 | } |
| 2889 | } |
| 2890 | |
| 2891 | fn parse_json_body<T: DeserializeOwned>(body: &[u8]) -> Result<T> { |
| 2892 | serde_json::from_slice(body).context("The Codewhale service returned an invalid JSON response") |
| 2893 | } |
| 2894 | |
| 2895 | /// A non-success reply from the account API, kept typed so a caller can tell a |
| 2896 | /// definitive refusal (4xx) from an outcome it cannot know (5xx, timeout). |
| 2897 | #[derive(Debug)] |
| 2898 | pub(crate) struct CloudHttpError { |
| 2899 | status: u16, |
| 2900 | code: Option<String>, |
| 2901 | reconciliation_required: bool, |
| 2902 | } |
| 2903 | |
| 2904 | impl CloudHttpError { |
| 2905 | fn code(&self) -> Option<&str> { |
| 2906 | self.code.as_deref() |
| 2907 | } |
| 2908 | } |
| 2909 | |
| 2910 | impl std::fmt::Display for CloudHttpError { |
| 2911 | fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result { |
| 2912 | match &self.code { |
| 2913 | Some(code) => write!( |
| 2914 | f, |
| 2915 | "Codewhale account request failed (HTTP {}, code {code})", |
| 2916 | self.status |
| 2917 | ), |
| 2918 | None => write!(f, "Codewhale account request failed (HTTP {})", self.status), |
| 2919 | } |
| 2920 | } |
| 2921 | } |
| 2922 | |
| 2923 | impl std::error::Error for CloudHttpError {} |
| 2924 | |
| 2925 | /// The request may or may not have been processed: it never reached the |
| 2926 | /// service, or its reply was lost or unreadable. Mutating commands use this to |
| 2927 | /// say "unknown" instead of guessing "failed". |
| 2928 | #[derive(Debug)] |
| 2929 | pub(crate) struct CloudTransportError { |
| 2930 | message: &'static str, |
| 2931 | source: Box<dyn std::error::Error + Send + Sync + 'static>, |
| 2932 | } |
| 2933 | |
| 2934 | impl CloudTransportError { |
| 2935 | fn new(message: &'static str, source: impl std::error::Error + Send + Sync + 'static) -> Self { |
| 2936 | Self { |
| 2937 | message, |
| 2938 | source: Box::new(source), |
| 2939 | } |
| 2940 | } |
| 2941 | } |
| 2942 | |
| 2943 | impl std::fmt::Display for CloudTransportError { |
| 2944 | fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result { |
| 2945 | f.write_str(self.message) |
| 2946 | } |
| 2947 | } |
| 2948 | |
| 2949 | impl std::error::Error for CloudTransportError { |
| 2950 | fn source(&self) -> Option<&(dyn std::error::Error + 'static)> { |
| 2951 | Some(self.source.as_ref()) |
| 2952 | } |
| 2953 | } |
| 2954 | |
| 2955 | /// Whether a failed mutating request may still have taken effect. |
| 2956 | fn outcome_unknown(err: &anyhow::Error) -> bool { |
| 2957 | err.downcast_ref::<CloudTransportError>().is_some() |
| 2958 | || err.downcast_ref::<CloudHttpError>().is_some_and(|http| { |
| 2959 | http.status >= 500 |
| 2960 | || http.status == 408 |
| 2961 | || http.reconciliation_required |
| 2962 | || http.code().is_some_and(|code| { |
| 2963 | code.ends_with("_outcome_unknown") |
| 2964 | || matches!( |
| 2965 | code, |
| 2966 | "boat_task_replay_expired" |
| 2967 | | "boat_task_create_in_progress" |
| 2968 | | "boat_task_cleanup_pending" |
| 2969 | | "boat_task_receipt_invalid" |
| 2970 | | "boat_task_authority_changed" |
| 2971 | | "boat_task_stop_unconfirmed" |
| 2972 | | "boat_task_usage_pending" |
| 2973 | ) |
| 2974 | }) |
| 2975 | }) |
| 2976 | } |
| 2977 | |
| 2978 | fn response_error(response: &CloudResponse) -> anyhow::Error { |
| 2979 | let body = serde_json::from_slice::<serde_json::Value>(&response.body).ok(); |
| 2980 | let code = body.as_ref().and_then(|body| { |
| 2981 | body.get("code") |
| 2982 | .and_then(serde_json::Value::as_str) |
| 2983 | .or_else(|| { |
| 2984 | body.get("error") |
| 2985 | .and_then(|error| error.get("code")) |
| 2986 | .and_then(serde_json::Value::as_str) |
| 2987 | }) |
| 2988 | .and_then(safe_error_code) |
| 2989 | }); |
| 2990 | let reconciliation_required = body.as_ref().is_some_and(|body| { |
| 2991 | body.get("reconciliationRequired") |
| 2992 | .or_else(|| { |
| 2993 | body.get("error") |
| 2994 | .and_then(|error| error.get("reconciliationRequired")) |
| 2995 | }) |
| 2996 | .and_then(serde_json::Value::as_bool) |
| 2997 | == Some(true) |
| 2998 | }); |
| 2999 | anyhow::Error::new(CloudHttpError { |
| 3000 | status: response.status, |
| 3001 | code, |
| 3002 | reconciliation_required, |
| 3003 | }) |
| 3004 | } |
| 3005 | |
| 3006 | fn safe_error_code(code: &str) -> Option<String> { |
| 3007 | if code.is_empty() |
| 3008 | || code.len() > 80 |
| 3009 | || !code |
| 3010 | .bytes() |
| 3011 | .all(|byte| byte.is_ascii_alphanumeric() || matches!(byte, b'_' | b'-' | b'.')) |
| 3012 | { |
| 3013 | return None; |
| 3014 | } |
| 3015 | Some(code.to_string()) |
| 3016 | } |
| 3017 | |
| 3018 | fn printable(value: &str) -> String { |
| 3019 | printable_max(value, 200) |
| 3020 | } |
| 3021 | |
| 3022 | /// `printable` with a caller-chosen bound, for remote prose longer than a label. |
| 3023 | fn printable_max(value: &str, max_chars: usize) -> String { |
| 3024 | value |
| 3025 | .chars() |
| 3026 | .filter(|character| !character.is_control()) |
| 3027 | .take(max_chars) |
| 3028 | .collect::<String>() |
| 3029 | .trim() |
| 3030 | .to_string() |
| 3031 | } |
| 3032 | |
| 3033 | #[cfg(test)] |
| 3034 | mod tests; |
| 3035 |