ocre-cli 0.2.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, loads htmx's
//! WebSocket extension in `templates/layout.html` 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";

/// htmx's WebSocket extension, loaded by the layout. A page that loaded it
/// itself would race on boosted navigations: htmx processes `ws-connect` in
/// the swapped page before the script arrives, and never connects.
const WS_SCRIPT: &str =
    r#"<script src="https://unpkg.com/htmx-ext-ws@2.0.4/dist/ws.js" crossorigin="anonymous"></script>"#;
/// The layout line the extension goes after.
const HTMX_SCRIPT: &str = r#"<script src="https://unpkg.com/htmx.org@"#;

/// 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> {
    add_channel_arm(edits, &format!("\"{channel}\" => {{}}"), command)
}

/// Lets browsers listen to the channels of `prefix:<anything>` (one per
/// record, named by an unguessable id).
pub(crate) fn add_channel_prefix(edits: &mut Edits, prefix: &str, command: &str) -> Result<(), CliError> {
    add_channel_arm(edits, &format!("channel if channel.starts_with(\"{prefix}:\") => {{}}"), command)
}

/// The feature, the Durable Object, the route, and `arm` in the `match` of `connect`.
fn add_channel_arm(edits: &mut Edits, arm: &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)?);
    }
    if let Some(layout) = edits.read(LAYOUT)?
        && !layout.contains(WS_SCRIPT)
    {
        edits.update(LAYOUT, with_ws_script(&layout)?);
    }
    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(())
}

const LAYOUT: &str = "templates/layout.html";

/// `layout` with [`WS_SCRIPT`] on the line after htmx's script tag, or
/// before `</head>` when htmx comes from elsewhere.
fn with_ws_script(layout: &str) -> Result<String, CliError> {
    if let Some(start) = layout.find(HTMX_SCRIPT) {
        let line_end = layout[start..].find('\n').map_or(layout.len(), |end| start + end);
        let line_start = layout[..start].rfind('\n').map_or(0, |newline| newline + 1);
        let indent = &layout[line_start..start];
        let indent = &indent[..indent.len() - indent.trim_start().len()];
        return Ok(format!("{}\n{indent}{WS_SCRIPT}{}", &layout[..line_end], &layout[line_end..]));
    }
    let head = layout.find("</head>").ok_or_else(|| {
        CliError::new("templates/layout.html has no </head>")
            .hint(format!("add `{WS_SCRIPT}` to the <head> of your layout, after htmx, then run the generator again"))
    })?;
    Ok(format!("{}  {WS_SCRIPT}\n{}", &layout[..head], &layout[head..]))
}

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` (and `ocre g external_job ... --realtime`) adds its channels 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
}}
"#
    )
}