File size: 4,043 Bytes
4d393e8
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
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
"""
`dw add` command.
Ingests local documents into DocWeave staging via upload_document.
"""
from __future__ import annotations

import os
from pathlib import Path
from typing import Optional
import typer
from rich.progress import Progress, SpinnerColumn, TextColumn

from diffweave.cli.config import get_active_workspace
from diffweave.cli.output import print_success, print_error, print_info, print_json
from diffweave.mcp.client import DiffWeaveMCPClient


def add_command(
    file_path: str = typer.Argument(..., help="Path to document file to stage and ingest"),
    workspace_id: Optional[str] = typer.Option(None, "--workspace", "-w", help="Workspace UUID override"),
    wait: bool = typer.Option(True, "--wait/--no-wait", help="Wait for extraction and proposal workflow to finish"),
    as_json: bool = typer.Option(False, "--json", help="Output machine-readable JSON"),
):
    """
    Stage and ingest a document into the workspace knowledge pipeline.
    """
    path = Path(file_path).resolve()
    if not path.exists():
        print_error(f"File not found: {file_path}")
        raise typer.Exit(code=1)

    try:
        active_ws_id = get_active_workspace(workspace_id)
    except Exception as e:
        print_error(str(e))
        raise typer.Exit(code=1)

    client = DiffWeaveMCPClient()

    if not as_json:
        print_info(f"Ingesting [bold white]{path.name}[/bold white] into workspace [cyan]{active_ws_id}[/cyan]...")

    with Progress(
        SpinnerColumn(),
        TextColumn("[progress.description]{task.description}"),
        transient=True,
    ) as progress:
        task = progress.add_task(f"Uploading {path.name} to DocWeave staging...", total=None)
        try:
            result = client.call_tool_sync(
                "upload_document",
                {"workspace_id": active_ws_id, "file_path": str(path)},
            )
        except Exception as e:
            progress.stop()
            print_error(f"Upload failed: {e}")
            raise typer.Exit(code=1)

    doc_id = result.get("document_id")
    wf_id = result.get("workflow_id")
    version_id = result.get("version_id")

    if wait and wf_id:
        import time
        with Progress(
            SpinnerColumn(),
            TextColumn("[progress.description]{task.description}"),
            transient=True,
        ) as progress:
            wf_task = progress.add_task("Running extraction, reconciliation & fact proposals...", total=None)
            for _ in range(60):  # Wait up to 120s
                time.sleep(2)
                try:
                    wf_status = client.call_tool_sync("get_workflow_status", {"workflow_id": wf_id})
                    curr_state = wf_status.get("status")
                    if curr_state in ("WAITING_FOR_REVIEW", "COMPLETED", "FAILED"):
                        break
                except Exception:
                    pass

    if as_json:
        print_json(result)
        return

    print_success(f"Document [bold]{path.name}[/bold] staged successfully!")
    print_info(f"  Document ID:  [dim]{doc_id}[/dim]")
    print_info(f"  Version ID:   [dim]{version_id}[/dim]")
    print_info(f"  Workflow ID:  [cyan]{wf_id}[/cyan]")

    # Check workflow final status
    try:
        wf_status = client.call_tool_sync("get_workflow_status", {"workflow_id": wf_id})
        status_val = wf_status.get("status", "RUNNING")
        if status_val == "WAITING_FOR_REVIEW":
            print_info(f"  Workflow State: [bold yellow]WAITING FOR REVIEW[/bold yellow] (Run 'dw diff' or 'dw review' to inspect)")
        elif status_val == "COMPLETED":
            print_success(f"  Workflow State: [bold green]COMPLETED[/bold green] (Knowledge items auto-approved)")
        elif status_val == "FAILED":
            print_error(f"  Workflow State: [bold red]FAILED[/bold red] - {wf_status.get('error_message', 'Unknown error')}")
        else:
            print_info(f"  Workflow State: [bold cyan]{status_val}[/bold cyan] (Run 'dw status' to track)")
    except Exception:
        pass