loopflow 0.12.3

Run steps and flows with coding agents
Documentation
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
217
218
219
220
221
222
223
224
225
226
227
228
229
230
231
232
233
234
235
236
237
238
239
240
241
242
243
244
245
246
247
248
249
250
251
252
253
254
255
256
257
258
259
260
261
262
263
264
265
266
267
268
269
270
271
272
273
274
275
276
277
278
279
280
281
282
283
284
285
286
287
288
289
290
291
292
293
294
295
296
297
298
299
300
301
302
303
304
305
306
307
308
309
310
311
312
313
314
315
316
317
318
319
320
321
322
323
324
325
326
327
328
329
330
331
332
333
334
335
336
337
338
339
340
341
342
343
344
345
346
347
348
349
350
351
352
353
354
355
356
357
358
359
360
361
362
363
364
365
366
367
368
369
370
371
372
373
374
375
376
377
378
379
380
381
382
383
384
385
386
387
388
389
390
391
392
393
394
395
396
397
398
399
400
401
402
403
404
405
406
407
408
409
410
411
412
413
414
415
416
417
418
419
420
421
422
423
424
425
426
427
428
429
430
431
432
433
434
435
436
437
438
439
//! `lf memory` — read or curate a wave's MEMORY.md through its live server.
//!
//! The live server holds the pen: `update` (full replacement from stdin) POSTs
//! the compiled checkpoint and `add "fact"` publishes one replayable fact to
//! the stream. `show` (the bare default) reads through the server when one is
//! live and falls back to the origin file otherwise. `log` prints add-stream
//! facts oldest to newest, using the live replay buffer or the journal fold.
//! Targeting (`--wave`, `--parent`) matches `lf chat` ([`super::chat`]),
//! including the drop rule: a write with no wave context anywhere is a publish
//! to no subscriber — exit 0, one stderr note. Reads require a wave context.

use std::io::Read;

use anyhow::{anyhow, Result};

use crate::lf::commands::chat::{get_json, post_json, resolve_target, CliContext, ResolvedWave};
use crate::lf::{MemoryCommand, WaveTargetArgs};
use crate::receipt::Receipt;
use crate::wave::journal::{fold_thread, journal_path, memory_facts, read_events, MemoryFact};
use crate::wave::memory::Memory;

pub fn run(cmd: Option<&MemoryCommand>, default_target: &WaveTargetArgs) -> Result<()> {
    let rt = tokio::runtime::Runtime::new()?;
    rt.block_on(async {
        let context = CliContext::detect().await;
        run_with_context(&context, cmd, default_target).await
    })
}

pub(crate) async fn run_with_context(
    context: &CliContext,
    cmd: Option<&MemoryCommand>,
    default_target: &WaveTargetArgs,
) -> Result<()> {
    match cmd {
        None => show(context, default_target).await,
        Some(MemoryCommand::Show { target }) => show(context, target).await,
        Some(MemoryCommand::Log { json, target }) => log(context, target, *json).await,
        Some(MemoryCommand::Update { summary, target }) => {
            // Resolve before touching stdin so a no-wave drop never blocks.
            let Some(resolved) = resolve(context, target).await? else {
                drop_note();
                return Ok(());
            };
            let mut content = String::new();
            std::io::stdin().read_to_string(&mut content)?;
            let summary =
                write_memory(&resolved, "update", &content, summary.as_deref(), &[]).await?;
            println!("memory updated for wave '{}': {summary}", resolved.name);
            Ok(())
        }
        Some(MemoryCommand::Add {
            fact,
            receipts,
            target,
        }) => {
            let Some(resolved) = resolve(context, target).await? else {
                drop_note();
                return Ok(());
            };
            // Parse at the CLI boundary so a bad `--receipt` is a user error the
            // author sees, and stamp each with the claim's own wave.
            let parsed = receipts
                .iter()
                .map(|token| Receipt::parse(token, &resolved.name))
                .collect::<std::result::Result<Vec<_>, _>>()
                .map_err(|err| anyhow!("{err}"))?;
            let summary = write_memory(&resolved, "add", fact, None, &parsed).await?;
            println!("memory fact added for wave '{}': {summary}", resolved.name);
            Ok(())
        }
    }
}

/// Publish-to-no-subscriber: writes outside any wave drop with exit 0.
fn drop_note() {
    eprintln!("no wave here; memory write dropped");
}

/// Memory is wave-level (MEMORY.md is wave identity; work lines have no
/// memory — their notes are files), so this resolves the FAMILY HEAD even
/// inside a work-line worktree and ignores the channel arm entirely.
async fn resolve(context: &CliContext, target: &WaveTargetArgs) -> Result<Option<ResolvedWave>> {
    resolve_target(
        target,
        context.store.as_ref(),
        context.repo.as_deref(),
        context.env_wave_id.as_deref(),
        context.env_channel.as_deref(),
    )
    .await
}

/// Reads are not publishes: no wave context is an error, not a drop.
async fn show(context: &CliContext, target: &WaveTargetArgs) -> Result<()> {
    let resolved = resolve(context, target).await?.ok_or_else(|| {
        anyhow!(
            "cannot resolve a target wave: no LF_WAVE_ID in env and \
             not inside a wave worktree — pass --wave <name>"
        )
    })?;
    print!("{}", read_memory(&resolved).await?);
    Ok(())
}

/// Reads are not publishes: no wave context is an error, not a drop.
async fn log(context: &CliContext, target: &WaveTargetArgs, json: bool) -> Result<()> {
    let resolved = resolve(context, target).await?.ok_or_else(|| {
        anyhow!(
            "cannot resolve a target wave: no LF_WAVE_ID in env and \
             not inside a wave worktree — pass --wave <name>"
        )
    })?;
    if json {
        let facts = read_memory_facts(&resolved)?;
        println!("{}", serde_json::to_string_pretty(&facts)?);
        return Ok(());
    }
    for fact in read_memory_log(&resolved).await? {
        println!("{fact}");
    }
    Ok(())
}

/// The wave's MEMORY.md: through the live server when one answers, else a
/// direct read of the origin file (reads don't need the pen).
pub(crate) async fn read_memory(resolved: &ResolvedWave) -> Result<String> {
    if let Some(endpoint) = &resolved.endpoint {
        let body = get_json(endpoint, "/memory").await?;
        return Ok(body["content"].as_str().unwrap_or_default().to_string());
    }
    let root = resolved.repo_root.as_deref().ok_or_else(|| {
        anyhow!(
            "wave '{}' has no live server and no local wave directory to read",
            resolved.name
        )
    })?;
    Ok(Memory::for_wave(root, &resolved.name).read())
}

/// Facts added to the wave's memory stream, oldest to newest. Prefer the live
/// server's replay buffer; fall back to the journal fold when no server is
/// running.
pub(crate) async fn read_memory_log(resolved: &ResolvedWave) -> Result<Vec<String>> {
    if let Some(endpoint) = &resolved.endpoint {
        let body = get_json(endpoint, "/memory/log").await?;
        return Ok(body["facts"]
            .as_array()
            .map(|facts| {
                facts
                    .iter()
                    .filter_map(|fact| fact.as_str().map(str::to_string))
                    .collect()
            })
            .unwrap_or_default());
    }
    let root = resolved.repo_root.as_deref().ok_or_else(|| {
        anyhow!(
            "wave '{}' has no live server and no local wave directory to read",
            resolved.name
        )
    })?;
    let events = read_events(&journal_path(root, &resolved.name));
    Ok(fold_thread(&events).memory_adds)
}

/// Facts with their evidence receipts, oldest to newest — the receipt-bearing
/// view behind `lf memory log --json`. Receipts live only in the journal (the
/// live `/memory/log` stream carries prose alone), so this reads the origin
/// journal directly. A live server appends to that same file under its lock, so
/// the read is current whether or not a server is running.
pub(crate) fn read_memory_facts(resolved: &ResolvedWave) -> Result<Vec<MemoryFact>> {
    let root = resolved.repo_root.as_deref().ok_or_else(|| {
        anyhow!(
            "wave '{}' has no local wave directory; receipts are journaled locally, \
             so `--json` needs the wave's worktree",
            resolved.name
        )
    })?;
    let events = read_events(&journal_path(root, &resolved.name));
    Ok(memory_facts(&events))
}

/// Write through the live server (the sole holder of MEMORY.md's pen).
/// Returns the summary the server journaled. No live server is an error —
/// there is deliberately no offline write path.
pub(crate) async fn write_memory(
    resolved: &ResolvedWave,
    op: &str,
    content: &str,
    summary: Option<&str>,
    receipts: &[Receipt],
) -> Result<String> {
    let endpoint = resolved.require_endpoint()?;
    let body = post_json(
        &endpoint,
        "/memory",
        &serde_json::json!({
            "op": op,
            "content": content,
            "summary": summary,
            "receipts": receipts,
        }),
    )
    .await?;
    Ok(body["summary"].as_str().unwrap_or_default().to_string())
}

#[cfg(test)]
mod tests {
    use super::*;
    use std::path::Path;

    use crate::lf::commands::fixtures::boot_server;
    use crate::wave::journal::{journal_path, EventKind, Journal};
    use crate::wave::runtime::WaveRuntime;

    fn resolved(name: &str, endpoint: Option<String>, root: Option<&Path>) -> ResolvedWave {
        ResolvedWave {
            name: name.to_string(),
            endpoint,
            repo_root: root.map(Path::to_path_buf),
        }
    }

    fn memory_events(origin: &Path, wave: &str) -> Vec<EventKind> {
        let (_, events) = Journal::open(&journal_path(origin, wave)).expect("journal");
        events.into_iter().map(|event| event.kind).collect()
    }

    /// `update` replaces the ORIGIN repo's file through the server and
    /// journals `MemoryUpdated`; `add` publishes a replayable fact without
    /// mutating the compiled file.
    #[tokio::test]
    async fn update_writes_the_origin_file_and_add_journals() {
        let tmp = tempfile::tempdir().expect("tempdir");
        let origin = tmp.path();
        let (addr, _runtime, _inbox) = boot_server(origin, "ship").await;
        let target = resolved("ship", Some(addr), None);

        let summary = write_memory(&target, "update", "# Ship\n\nfold is truth\n", None, &[])
            .await
            .expect("update");
        assert_eq!(summary, "# Ship", "summary defaults to the first line");
        assert_eq!(
            std::fs::read_to_string(origin.join("wave/ship/MEMORY.md")).expect("origin file"),
            "# Ship\n\nfold is truth\n",
            "the ORIGIN file is the one replaced"
        );

        let summary = write_memory(&target, "add", "workers report via lf radio pub", None, &[])
            .await
            .expect("add");
        assert_eq!(summary, "workers report via lf radio pub");
        assert_eq!(
            std::fs::read_to_string(origin.join("wave/ship/MEMORY.md")).expect("origin file"),
            "# Ship\n\nfold is truth\n",
            "add publishes a stream fact without accreting raw bullets"
        );
        assert_eq!(
            memory_events(origin, "ship")
                .into_iter()
                .filter(|event| matches!(
                    event,
                    EventKind::MemoryUpdated { .. } | EventKind::MemoryAdded { .. }
                ))
                .collect::<Vec<_>>(),
            vec![
                EventKind::MemoryUpdated {
                    summary: "# Ship".to_string()
                },
                EventKind::MemoryAdded {
                    fact: "workers report via lf radio pub".to_string(),
                    receipts: Vec::new(),
                },
            ]
        );
    }

    #[tokio::test]
    async fn log_reads_through_the_server_when_live() {
        let tmp = tempfile::tempdir().expect("tempdir");
        let origin = tmp.path();
        let (addr, runtime, _inbox) = boot_server(origin, "ship").await;
        runtime.append_memory("first", vec![]).expect("append");
        runtime.append_memory("second", vec![]).expect("append");

        let facts = read_memory_log(&resolved("ship", Some(addr), None))
            .await
            .expect("read log");
        assert_eq!(facts, vec!["first".to_string(), "second".to_string()]);
    }

    #[tokio::test]
    async fn log_reads_the_journal_without_a_server() {
        let tmp = tempfile::tempdir().expect("tempdir");
        let origin = tmp.path();
        {
            let runtime =
                WaveRuntime::open("ship".to_string(), origin.to_path_buf()).expect("runtime");
            runtime
                .append_memory("offline first", vec![])
                .expect("append");
            runtime
                .append_memory("offline second", vec![])
                .expect("append");
            runtime
                .update_memory("# Ship\n\ncompiled\n", "compiled")
                .expect("update");
            runtime
                .append_memory("after update", vec![])
                .expect("append");
        }

        let facts = read_memory_log(&resolved("ship", None, Some(origin)))
            .await
            .expect("read log");
        assert_eq!(facts, vec!["after update".to_string()]);
    }

    /// End to end: `add --receipt` writes through the live server, the receipt is
    /// journaled with the fact, and `read_memory_facts` reads it back — the
    /// authoring-and-drill path the whole task rests on.
    #[tokio::test]
    async fn add_writes_receipts_that_the_json_view_reads_back() {
        use crate::receipt::{EvidenceKind, Receipt};

        let tmp = tempfile::tempdir().expect("tempdir");
        let origin = tmp.path();
        let (addr, _runtime, _inbox) = boot_server(origin, "ship").await;
        // A local wave is both served and on disk: the receipt view reads the
        // journal even while the server holds the pen.
        let target = resolved("ship", Some(addr), Some(origin));

        let receipts = vec![
            Receipt::new(EvidenceKind::ChatTurn, "turn-3", "ship"),
            Receipt::new(EvidenceKind::Pr, "loopflow/loopflow#912@abc1234", "ship"),
        ];
        write_memory(
            &target,
            "add",
            "workers report via the stream",
            None,
            &receipts,
        )
        .await
        .expect("add with receipts");

        let facts = read_memory_facts(&target).expect("read facts");
        assert_eq!(facts.len(), 1);
        assert_eq!(facts[0].fact, "workers report via the stream");
        assert_eq!(facts[0].receipts, receipts);
    }

    /// `show` works with no server at all: a direct read of the origin file.
    #[tokio::test]
    async fn show_reads_the_origin_file_without_a_server() {
        let tmp = tempfile::tempdir().expect("tempdir");
        let dir = tmp.path().join("wave/ship");
        std::fs::create_dir_all(&dir).unwrap();
        std::fs::write(dir.join("MEMORY.md"), "offline read\n").unwrap();

        let content = read_memory(&resolved("ship", None, Some(tmp.path())))
            .await
            .expect("read");
        assert_eq!(content, "offline read\n");
    }

    /// `show` prefers the live server's view when one answers.
    #[tokio::test]
    async fn show_reads_through_the_server_when_live() {
        let tmp = tempfile::tempdir().expect("tempdir");
        let dir = tmp.path().join("wave/ship");
        std::fs::create_dir_all(&dir).unwrap();
        std::fs::write(dir.join("MEMORY.md"), "served content\n").unwrap();
        let (addr, _runtime, _inbox) = boot_server(tmp.path(), "ship").await;

        let content = read_memory(&resolved("ship", Some(addr), None))
            .await
            .expect("read");
        assert_eq!(content, "served content\n");
    }

    /// No wave context at all: a write is a publish to no subscriber — it
    /// drops with exit 0 instead of erroring.
    #[tokio::test]
    async fn add_without_wave_context_drops_with_exit_zero() {
        let tmp = tempfile::tempdir().expect("tempdir");
        let context = CliContext {
            store: None,
            repo: Some(tmp.path().to_path_buf()),
            env_wave_id: None,
            env_channel: None,
        };
        run_with_context(
            &context,
            Some(&MemoryCommand::Add {
                fact: "dropped fact".to_string(),
                receipts: Vec::new(),
                target: WaveTargetArgs::default(),
            }),
            &WaveTargetArgs::default(),
        )
        .await
        .expect("dropped write exits 0");
    }

    /// `show` is a read, not a publish: no wave context stays a clear error.
    #[tokio::test]
    async fn show_without_wave_context_errors() {
        let tmp = tempfile::tempdir().expect("tempdir");
        let context = CliContext {
            store: None,
            repo: Some(tmp.path().to_path_buf()),
            env_wave_id: None,
            env_channel: None,
        };
        let err = run_with_context(&context, None, &WaveTargetArgs::default())
            .await
            .expect_err("read with no wave context");
        assert!(err.to_string().contains("--wave"), "{err}");
    }

    /// Writes have no offline path: no live server is an error.
    #[tokio::test]
    async fn update_without_a_server_errors() {
        let tmp = tempfile::tempdir().expect("tempdir");
        let err = write_memory(
            &resolved("ship", None, Some(tmp.path())),
            "update",
            "x",
            None,
            &[],
        )
        .await
        .expect_err("no server");
        assert!(err.to_string().contains("no live listener"), "{err}");
    }
}