ocre-cli 0.1.0

Command-line tool for Ocre: create, generate, migrate, run and deploy apps.
//! The app side of `ocre::realtime`, set up by `ocre g scaffold ... --realtime`.
//! The first use turns on Ocre's `realtime` feature, binds and exports the
//! `OcreChannel` Durable Object in cloudflare.config.ts and writes
//! `src/realtime.rs` (the `/realtime/{channel}` route and who may listen);
//! later uses add their channel to it.

use super::{Edits, insert_after_marker, read_config, register_routes, with_ocre_feature};
use crate::{
    config::{self, ENV_MARKER, EXPORTS_MARKER},
    output::CliError,
};

const CHANNELS_MARKER: &str = "// ocre:channels";

/// Exported once. Deploying creates the namespace from the export, so
/// `ocre deploy` has nothing to create beforehand.
const CHANNELS_EXPORT: &str = "// Realtime channels (`ocre::realtime`): one Durable Object per channel holds the
// browsers' WebSockets, hibernated between broadcasts so idle connections cost
// no duration. The free plan only accepts SQLite-backed classes.
OcreChannel: exports.durableObject({ storage: \"sqlite\" }),";

/// Lets browsers listen to `channel`: feature, Durable Object and route.
pub(crate) fn add_channel(edits: &mut Edits, channel: &str, command: &str) -> Result<(), CliError> {
    let cargo = edits.read("Cargo.toml")?.unwrap_or_default();
    let with_feature = with_ocre_feature(&cargo, "realtime")?;
    if with_feature != cargo {
        edits.update("Cargo.toml", with_feature);
    }
    let config = read_config(edits)?;
    if config.export("OcreChannel").is_none() {
        edits.update(config::FILE, config.insert(EXPORTS_MARKER, CHANNELS_EXPORT)?);
    }
    let config = read_config(edits)?;
    if config.binding("CHANNELS").is_none() {
        let binding = format!(
            "// The realtime channels' Durable Objects (`ocre::realtime`).\nCHANNELS: bindings.durableObject({{ worker: \"{}\", exportName: \"OcreChannel\" }}),",
            config.worker_name()?
        );
        edits.update(config::FILE, config.insert(ENV_MARKER, &binding)?);
    }
    let arm = format!("\"{channel}\" => {{}}");
    match edits.read("src/realtime.rs")? {
        Some(module) if module.contains(&arm) => {}
        Some(module) => {
            let module = insert_after_marker(&module, CHANNELS_MARKER, &arm).ok_or_else(|| {
                CliError::new("src/realtime.rs is missing the `// ocre:channels` marker").hint(format!(
                    "put `// ocre:channels` on its own line in the `match channel.as_str()` of `connect`, or add `{arm}` there yourself"
                ))
            })?;
            edits.update("src/realtime.rs", module);
        }
        None => {
            edits.create("src/realtime.rs", module_rs(&arm, command))?;
            register_routes(edits, "realtime")?;
        }
    }
    Ok(())
}

fn module_rs(arm: &str, command: &str) -> String {
    format!(
        r#"//! Realtime channels: pages subscribe with htmx's WebSocket extension
//! (`<div hx-ext="ws" ws-connect="/realtime/<channel>">`) and receive the HTML
//! that handlers and jobs send with `ocre::realtime::broadcast(&ctx, "<channel>", &html)`.
//! Generated by `{command}`.
//! Each `ocre g scaffold ... --realtime` adds its channel below.

use axum::{{
    Router,
    extract::{{Path, State}},
    response::Response,
    routing::get,
}};
use ocre::{{Ctx, Error, Result, realtime::WebSocketUpgrade}};

/// `/ocre/dev/realtime/sent.json` lists recent broadcasts in `ocre dev` (for
/// tests); it is a 404 in deployed builds.
pub fn routes() -> Router<Ctx> {{
    Router::new().route("/realtime/{{channel}}", get(connect)).merge(ocre::realtime::dev_routes())
}}

/// Opens a WebSocket on `channel` for whoever may listen to it. Everyone may
/// listen to the channels below. To restrict one, extract the user here
/// (`crate::auth::CurrentUser` after `ocre g auth`) and return
/// `Err(Error::Forbidden)`; `upgrade.identified_by(user.id.to_string())`
/// names the subscriber. Unknown channels are a 404.
async fn connect(State(ctx): State<Ctx>, Path(channel): Path<String>, upgrade: WebSocketUpgrade) -> Result<Response> {{
    match channel.as_str() {{
        {CHANNELS_MARKER}
        {arm}
        _ => return Err(Error::NotFound),
    }}
    upgrade.connect(&ctx, &channel).await
}}
"#
    )
}