AkurAI Build
Menu

akurai-tasks

public

Latest change 4cb9b5f41b640c25c028037e0f48dac6d3956ad0 - fix(mcp): constrain work query schema by Ólafur Búi Ólafsson

#![forbid(unsafe_code)]

use akurai_json::{parse, Value};
use tasks_core::{
    field, object, string, CreateItem, CreateOutcome, LeaseProof, MetricsQuery, OutcomePatch,
    Query, Store, TransitionWebhook, STATES,
};

#[derive(Clone)]
pub struct Mcp {
    store: Store,
    actor: String,
    webhook: Option<TransitionWebhook>,
}

impl Mcp {
    pub fn new(store: Store, actor: String) -> Self {
        Self {
            store,
            actor,
            webhook: TransitionWebhook::from_env(),
        }
    }

    pub fn handle_line(&self, line: &str) -> String {
        let request = match parse(line) {
            Ok(request) => request,
            Err(error) => return rpc_error(Value::Null, -32700, &format!("parse error: {error}")),
        };
        let id = request.get("id").cloned();
        if id.is_none() {
            return String::new();
        }
        let id = id.unwrap_or(Value::Null);
        let method = field(&request, "method").unwrap_or("");
        let result = match method {
            "initialize" => Ok(object(vec![
                ("protocolVersion", string("2025-11-25")),
                (
                    "capabilities",
                    object(vec![
                        ("tools", object(vec![])),
                        ("resources", object(vec![])),
                        ("prompts", object(vec![])),
                    ]),
                ),
                (
                    "serverInfo",
                    object(vec![
                        ("name", string("akurai-tasks")),
                        ("version", string(env!("CARGO_PKG_VERSION"))),
                    ]),
                ),
            ])),
            "ping" => Ok(object(vec![])),
            "tools/list" => Ok(object(vec![("tools", Value::Array(tools()))])),
            "tools/call" => self.call(request.get("params").unwrap_or(&Value::Null)),
            "resources/list" => Ok(object(vec![(
                "resources",
                Value::Array(vec![object(vec![
                    ("uri", string("tasks://projects")),
                    ("name", string("Projects")),
                    ("mimeType", string("application/json")),
                ])]),
            )])),
            "resources/read" => self.read_resource(request.get("params").unwrap_or(&Value::Null)),
            "prompts/list" => Ok(object(vec![(
                "prompts",
                Value::Array(vec![object(vec![
                    ("name", string("next-task")),
                    (
                        "description",
                        string("Select the ready execution frontier for a project"),
                    ),
                    (
                        "arguments",
                        Value::Array(vec![object(vec![
                            ("name", string("project")),
                            ("required", Value::Bool(true)),
                        ])]),
                    ),
                ])]),
            )])),
            "prompts/get" => self.get_prompt(request.get("params").unwrap_or(&Value::Null)),
            _ => return rpc_error(id, -32601, &format!("method not found: {method}")),
        };
        match result {
            Ok(value) => rpc_result(id, value),
            Err(message) => rpc_error(id, -32602, &message),
        }
    }

    fn call(&self, params: &Value) -> Result<Value, String> {
        let name = required(params, "name")?;
        let arguments = params.get("arguments").unwrap_or(&Value::Null);
        let idempotency = field(arguments, "idempotencyKey").unwrap_or("");
        let result = match name {
            "project_list" | "projects_list" => {
                self.require_read_scope(&[])?;
                self.store.list_projects()
            }
            "project_create" => self.store.create_project(
                required(arguments, "key")?,
                required(arguments, "name")?,
                string_array(arguments, "repos")?,
                &self.actor,
                idempotency,
            ),
            "project_update" => self.store.update_project(
                required(arguments, "key")?,
                string_array(arguments, "repos")?,
                &self.actor,
                idempotency,
            ),
            "project_archive" => self.store.archive_project(
                required(arguments, "key")?,
                true,
                &self.actor,
                idempotency,
            ),
            "project_unarchive" => self.store.archive_project(
                required(arguments, "key")?,
                false,
                &self.actor,
                idempotency,
            ),
            "project_delete" => {
                self.store
                    .delete_project(required(arguments, "key")?, &self.actor, idempotency)
            }
            "project_purge" => {
                self.store
                    .purge_project(required(arguments, "key")?, &self.actor, idempotency)
            }
            "work_set_repo" | "task_set_repo" => self.store.set_item_repo(
                required(arguments, "id")?,
                Some(required(arguments, "repo")?),
                lease_proof(arguments)?,
                &self.actor,
                idempotency,
            ),
            "task_list" | "work_list" => {
                let project = required(arguments, "project")?.to_string();
                self.require_read_scope(std::slice::from_ref(&project))?;
                self.store.list_items(&project)
            }
            "task_show" | "work_get" => self
                .read_item(required(arguments, "id")?)
                .map_err(tasks_core::Error::Invalid),
            "task_create" | "work_create" => self.store.create_item(CreateItem {
                project: required(arguments, "project")?,
                title: required(arguments, "title")?,
                description: field(arguments, "description").unwrap_or(""),
                repo: field(arguments, "repo"),
                priority: field(arguments, "priority").unwrap_or("normal"),
                actor: &self.actor,
                idempotency_key: idempotency,
            }),
            "task_transition" | "work_transition" => {
                let id = required(arguments, "id")?;
                let previous = self.store.get_item(id).map_err(|error| error.to_string())?;
                let result = self
                    .store
                    .transition(
                        id,
                        required(arguments, "state")?,
                        integer(arguments, "expectedRevision")?,
                        lease_proof(arguments)?,
                        &self.actor,
                        idempotency,
                    )
                    .map_err(|error| error.to_string())?;
                if field(&previous, "state") != field(&result, "state")
                    && previous.get("revision").and_then(Value::as_i64)
                        < result.get("revision").and_then(Value::as_i64)
                {
                    if let Some(webhook) = &self.webhook {
                        webhook.enqueue(&previous, &result, &self.actor);
                    }
                }
                Ok(result)
            }
            "task_claim" | "work_claim" => self.store.claim(
                required(arguments, "id")?,
                &self.actor,
                integer(arguments, "expectedRevision")?,
                &self.actor,
                idempotency,
            ),
            "task_comment" | "task_evidence" => self.store.add_record(
                if name == "task_comment" {
                    "comment"
                } else {
                    "evidence"
                },
                required(arguments, "id")?,
                required(arguments, "body")?,
                lease_proof(arguments)?,
                &self.actor,
                idempotency,
            ),
            "task_dependency" => self.store.add_dependency(
                required(arguments, "id")?,
                required(arguments, "dependsOn")?,
                &self.actor,
                idempotency,
            ),
            "task_handoff" => self.store.handoff(
                required(arguments, "id")?,
                required(arguments, "to")?,
                required(arguments, "summary")?,
                lease_proof(arguments)?,
                &self.actor,
                idempotency,
            ),
            "task_accept_handoff" => {
                self.store
                    .accept_handoff(required(arguments, "id")?, &self.actor, idempotency)
            }
            "work_block" => self.store.set_blocked(
                required(arguments, "id")?,
                Some(required(arguments, "reason")?),
                integer(arguments, "expectedRevision")?,
                lease_proof(arguments)?,
                &self.actor,
                idempotency,
            ),
            "work_unblock" => self.store.set_blocked(
                required(arguments, "id")?,
                None,
                integer(arguments, "expectedRevision")?,
                lease_proof(arguments)?,
                &self.actor,
                idempotency,
            ),
            "outcome_create" => {
                let projects = string_array(arguments, "projects")?;
                let repos = string_array(arguments, "repos")?;
                let result = self
                    .store
                    .create_outcome(
                        CreateOutcome {
                            title: required(arguments, "title")?.to_string(),
                            description: field(arguments, "description").unwrap_or("").to_string(),
                            projects,
                            repos,
                            owner: required(arguments, "owner")?.to_string(),
                            target_metric: required(arguments, "targetMetric")?.to_string(),
                            baseline: arguments.get("baseline").cloned().unwrap_or(Value::Null),
                            target: arguments.get("target").cloned().unwrap_or(Value::Null),
                            observation_start: integer(arguments, "observationStart")?,
                            observation_end: integer(arguments, "observationEnd")?,
                            result_links: string_array(arguments, "resultLinks")?,
                        },
                        &self.actor,
                        idempotency,
                    )
                    .map_err(|error| error.to_string())?;
                self.authorize_outcome_read(&result)?;
                Ok(result)
            }
            "outcome_get" => {
                let result = self
                    .store
                    .get_outcome(required(arguments, "id")?)
                    .map_err(|error| error.to_string())?;
                self.authorize_outcome_read(&result)?;
                Ok(result)
            }
            "outcome_update" => {
                let id = required(arguments, "id")?;
                let patch = outcome_patch(arguments)?;
                let result = self
                    .store
                    .update_outcome(
                        id,
                        integer(arguments, "expectedRevision")?,
                        patch,
                        &self.actor,
                        idempotency,
                    )
                    .map_err(|error| error.to_string())?;
                self.authorize_outcome_read(&result)?;
                Ok(result)
            }
            "outcome_link" => {
                let result = self
                    .store
                    .link_outcome(
                        required(arguments, "outcome")?,
                        required(arguments, "item")?,
                        &self.actor,
                        idempotency,
                    )
                    .map_err(|error| error.to_string())?;
                Ok(result)
            }
            "outcome_unlink" => self.store.unlink_outcome(
                required(arguments, "outcome")?,
                required(arguments, "item")?,
                &self.actor,
                idempotency,
            ),
            "outcome_evidence" => self.store.add_outcome_evidence(
                required(arguments, "outcome")?,
                required(arguments, "state")?,
                required(arguments, "body")?,
                string_array(arguments, "resultLinks")?,
                &self.actor,
                idempotency,
            ),
            "outcome_query" | "portfolio_query" => {
                let projects = string_array(arguments, "projects")?;
                self.require_read_scope(&projects)?;
                self.store.query_outcomes(
                    projects,
                    string_array(arguments, "repos")?,
                    string_array(arguments, "owners")?,
                    string_array(arguments, "states")?,
                    arguments
                        .get("limit")
                        .and_then(Value::as_i64)
                        .unwrap_or(100)
                        .try_into()
                        .map_err(|_| "limit must be positive")?,
                )
            }
            "metrics_query" | "metrics_aggregate" => {
                let projects = string_array(arguments, "projects")?;
                self.require_read_scope(&projects)?;
                self.store.metrics(&MetricsQuery {
                    projects,
                    repos: string_array(arguments, "repos")?,
                    owners: string_array(arguments, "owners")?,
                    cycles: string_array(arguments, "cycles")?,
                    period: field(arguments, "period").unwrap_or("7d").to_string(),
                    as_of: arguments.get("asOf").and_then(Value::as_i64),
                })
            }
            "task_query" | "work_query" => {
                let query = query_from(arguments)?;
                self.require_read_scope(&query.projects)?;
                self.store.query_items(&query)
            }
            "work_renew" => self.store.renew_lease(
                required(arguments, "id")?,
                lease_proof(arguments)?.ok_or("lease token and generation are required")?,
                integer(arguments, "expectedRevision")?,
                &self.actor,
                idempotency,
            ),
            "work_release" => self.store.release_lease(
                required(arguments, "id")?,
                lease_proof(arguments)?.ok_or("lease token and generation are required")?,
                integer(arguments, "expectedRevision")?,
                &self.actor,
                idempotency,
            ),
            "import_apply" => self.store.apply_import(
                arguments.get("bundle").ok_or("bundle is required")?,
                &self.actor,
                idempotency,
            ),
            "import_verify" => {
                self.store
                    .verify_import(required(arguments, "run")?, &self.actor, idempotency)
            }
            "board_get" => {
                let project = required(arguments, "project")?.to_string();
                self.require_read_scope(std::slice::from_ref(&project))?;
                self.store.board(&project)
            }
            "ready_frontier" => {
                let project = required(arguments, "project")?.to_string();
                self.require_read_scope(std::slice::from_ref(&project))?;
                self.store.ready_frontier(&project)
            }
            "events_read" => {
                let project = field(arguments, "project").map(str::to_string);
                let projects = project.iter().cloned().collect::<Vec<_>>();
                self.require_read_scope(&projects)?;
                self.store.events(project.as_deref())
            }
            "system_doctor" => {
                self.require_admin()?;
                self.store.doctor()
            }
            "system_backup" => self.store.backup_snapshot(&self.actor),
            _ => return Err(format!("unknown tool: {name}")),
        }
        .map_err(|error| error.to_string())?;
        Ok(object(vec![(
            "content",
            Value::Array(vec![object(vec![
                ("type", string("text")),
                ("text", string(&result.to_json())),
            ])]),
        )]))
    }

    fn read_resource(&self, params: &Value) -> Result<Value, String> {
        match required(params, "uri")? {
            "tasks://projects" => {
                self.require_read_scope(&[])?;
                let projects = self
                    .store
                    .list_projects()
                    .map_err(|error| error.to_string())?;
                Ok(object(vec![(
                    "contents",
                    Value::Array(vec![object(vec![
                        ("uri", string("tasks://projects")),
                        ("mimeType", string("application/json")),
                        ("text", string(&projects.to_json())),
                    ])]),
                )]))
            }
            uri => Err(format!("unknown resource: {uri}")),
        }
    }

    fn get_prompt(&self, params: &Value) -> Result<Value, String> {
        if required(params, "name")? != "next-task" {
            return Err("unknown prompt".into());
        }
        let project = params
            .get("arguments")
            .and_then(|value| field(value, "project"))
            .ok_or_else(|| "project is required".to_string())?;
        Ok(object(vec![
            ("description", string("Pull the next policy-valid leaf task")),
            (
                "messages",
                Value::Array(vec![object(vec![
                    ("role", string("user")),
                    (
                        "content",
                        object(vec![
                            ("type", string("text")),
                            (
                                "text",
                                string(&format!(
                                    "Call ready_frontier for project {project}, select one leaf item, then claim it with its current revision."
                                )),
                            ),
                        ]),
                    ),
                ])]),
            ),
        ]))
    }
    fn read_item(&self, id: &str) -> Result<Value, String> {
        let item = self.store.get_item(id).map_err(|error| error.to_string())?;
        let project = field(&item, "project").ok_or("work item project is missing")?;
        self.require_read_scope(std::slice::from_ref(&project.to_string()))?;
        Ok(item)
    }

    fn require_admin(&self) -> Result<(), String> {
        if self
            .store
            .authorize(&self.actor, None, tasks_core::Permission::Admin)
            .map_err(|error| error.to_string())?
        {
            Ok(())
        } else {
            Err("administrator permission is required".into())
        }
    }

    fn require_read_scope(&self, projects: &[String]) -> Result<(), String> {
        if !self.store.bootstrap_complete().map_err(|e| e.to_string())? {
            return Ok(());
        }
        if projects.is_empty() {
            if self
                .store
                .authorize(&self.actor, None, tasks_core::Permission::Read)
                .map_err(|e| e.to_string())?
            {
                Ok(())
            } else {
                Err("cross-project read permission is required".into())
            }
        } else {
            for project in projects {
                if !self
                    .store
                    .authorize(&self.actor, Some(project), tasks_core::Permission::Read)
                    .map_err(|e| e.to_string())?
                {
                    return Err(format!(
                        "principal lacks read permission for project {project}"
                    ));
                }
            }
            Ok(())
        }
    }

    fn authorize_outcome_read(&self, outcome: &Value) -> Result<(), String> {
        let projects = string_array(outcome, "projects")?;
        self.require_read_scope(&projects)
    }
}

fn tools() -> Vec<Value> {
    [
        ("projects_list", "List all visible projects", &[][..]),
        (
            "project_create",
            "Create a project with one or more repositories",
            &["key", "name", "idempotencyKey"],
        ),
        (
            "project_update",
            "Replace the repositories registered to a project",
            &["key", "repos", "idempotencyKey"],
        ),
        (
            "project_archive",
            "Archive a project (reversible; hides it from active work)",
            &["key", "idempotencyKey"],
        ),
        (
            "project_unarchive",
            "Restore a previously archived project",
            &["key", "idempotencyKey"],
        ),
        (
            "project_delete",
            "Permanently delete an empty project (fails if it still has work items)",
            &["key", "idempotencyKey"],
        ),
        (
            "project_purge",
            "Permanently delete a project AND all its work items and history (irreversible)",
            &["key", "idempotencyKey"],
        ),
        (
            "work_set_repo",
            "Set the repository a work item belongs to",
            &["id", "repo", "idempotencyKey"],
        ),
        ("work_list", "List work items in one project", &["project"]),
        (
            "work_query",
            "Query work across projects and repositories using structured filters",
            &[][..],
        ),
        ("work_get", "Read one work item", &["id"]),
        (
            "work_create",
            "Create a work item",
            &["project", "title", "idempotencyKey"],
        ),
        (
            "work_transition",
            "Move a work item through its workflow",
            &["id", "state", "expectedRevision", "idempotencyKey"],
        ),
        (
            "work_block",
            "Block a work item with a reason",
            &["id", "reason", "expectedRevision", "idempotencyKey"],
        ),
        (
            "work_unblock",
            "Clear a work item's blocker",
            &["id", "expectedRevision", "idempotencyKey"],
        ),
        (
            "work_claim",
            "Claim a Ready work item and receive a fenced lease",
            &["id", "expectedRevision", "idempotencyKey"],
        ),
        (
            "work_renew",
            "Renew the current fenced lease",
            &[
                "id",
                "leaseToken",
                "leaseGeneration",
                "expectedRevision",
                "idempotencyKey",
            ],
        ),
        (
            "work_release",
            "Release the current fenced lease",
            &[
                "id",
                "leaseToken",
                "leaseGeneration",
                "expectedRevision",
                "idempotencyKey",
            ],
        ),
        (
            "task_comment",
            "Add a comment",
            &["id", "body", "idempotencyKey"],
        ),
        (
            "task_evidence",
            "Attach evidence",
            &["id", "body", "idempotencyKey"],
        ),
        (
            "task_dependency",
            "Add a work dependency",
            &["id", "dependsOn", "idempotencyKey"],
        ),
        (
            "task_handoff",
            "Request ownership handoff",
            &["id", "to", "summary", "idempotencyKey"],
        ),
        (
            "task_accept_handoff",
            "Accept ownership handoff",
            &["id", "idempotencyKey"],
        ),
        (
            "outcome_create",
            "Create a versioned strategic outcome spanning projects and repositories",
            &[
                "title",
                "projects",
                "owner",
                "targetMetric",
                "baseline",
                "target",
                "observationStart",
                "observationEnd",
                "idempotencyKey",
            ],
        ),
        ("outcome_get", "Read one strategic outcome", &["id"]),
        (
            "outcome_update",
            "Create a new audited outcome revision",
            &["id", "expectedRevision", "idempotencyKey"],
        ),
        (
            "outcome_link",
            "Link supporting work to an outcome",
            &["outcome", "item", "idempotencyKey"],
        ),
        (
            "outcome_unlink",
            "Remove supporting work from an outcome",
            &["outcome", "item", "idempotencyKey"],
        ),
        (
            "outcome_evidence",
            "Record explicit evidence that derives outcome state",
            &["outcome", "state", "body", "idempotencyKey"],
        ),
        (
            "outcome_query",
            "Query strategic outcomes across permitted projects",
            &[][..],
        ),
        (
            "portfolio_query",
            "Query strategic outcomes across permitted projects",
            &[][..],
        ),
        (
            "metrics_query",
            "Read bounded WIP, age and flow aggregates with source IDs",
            &["period"],
        ),
        ("board_get", "Read a project Kanban board", &["project"]),
        (
            "ready_frontier",
            "List dependency-free Ready work",
            &["project"],
        ),
        (
            "events_read",
            "Read immutable audit events, optionally scoped to a project",
            &[][..],
        ),
        (
            "system_doctor",
            "Check database integrity and report record counts",
            &[][..],
        ),
        (
            "system_backup",
            "Create a verified database snapshot in the configured staging directory",
            &[][..],
        ),
        (
            "import_apply",
            "Apply one lossless source import bundle",
            &["bundle", "idempotencyKey"],
        ),
        (
            "import_verify",
            "Verify imported records and exact-byte digests",
            &["run"],
        ),
    ]
    .into_iter()
    .map(|(name, description, required)| {
        let properties = if name == "work_query" {
            query_properties()
        } else if name == "outcome_query" || name == "portfolio_query" {
            outcome_query_properties()
        } else if name == "metrics_query" || name == "metrics_aggregate" {
            metrics_properties()
        } else {
            let mut fields: Vec<(&str, Value)> = required
                .iter()
                .map(|field| (*field, schema_for(field)))
                .collect();
            if matches!(
                name,
                "work_transition"
                    | "work_block"
                    | "work_unblock"
                    | "task_comment"
                    | "task_evidence"
                    | "task_handoff"
            ) {
                fields.push(("leaseToken", schema_for(&"leaseToken")));
                fields.push(("leaseGeneration", schema_for(&"leaseGeneration")));
            }
            object(fields)
        };
        object(vec![
            ("name", string(name)),
            ("description", string(description)),
            (
                "inputSchema",
                object(vec![
                    ("type", string("object")),
                    ("properties", properties),
                    (
                        "required",
                        Value::Array(required.iter().map(|name| string(name)).collect()),
                    ),
                    ("additionalProperties", Value::Bool(true)),
                ]),
            ),
        ])
    })
    .collect()
}

fn query_properties() -> Value {
    object(vec![
        ("projects", array_schema()),
        ("repos", array_schema()),
        ("states", state_array_schema()),
        ("priorities", array_schema()),
        ("owners", array_schema()),
        ("types", array_schema()),
        ("labels", array_schema()),
        ("parent", schema_for(&"parent")),
        ("cycle", schema_for(&"cycle")),
        ("text", schema_for(&"text")),
        ("blocked", object(vec![("type", string("boolean"))])),
        ("createdAfter", integer_schema()),
        ("createdBefore", integer_schema()),
        ("updatedAfter", integer_schema()),
        ("updatedBefore", integer_schema()),
        (
            "sort",
            object(vec![
                ("type", string("string")),
                (
                    "enum",
                    Value::Array(
                        ["id", "created", "updated", "priority", "state", "title"]
                            .iter()
                            .map(|v| string(v))
                            .collect(),
                    ),
                ),
            ]),
        ),
        ("descending", object(vec![("type", string("boolean"))])),
        ("limit", bounded_integer_schema(1, 500)),
        ("byteLimit", bounded_integer_schema(1024, 16 * 1024 * 1024)),
        ("cursor", schema_for(&"cursor")),
    ])
}

fn outcome_query_properties() -> Value {
    object(vec![
        ("projects", array_schema()),
        ("repos", array_schema()),
        ("owners", array_schema()),
        ("states", array_schema()),
        ("limit", integer_schema()),
    ])
}

fn metrics_properties() -> Value {
    object(vec![
        ("projects", array_schema()),
        ("repos", array_schema()),
        ("owners", array_schema()),
        ("cycles", array_schema()),
        (
            "period",
            object(vec![
                ("type", string("string")),
                (
                    "enum",
                    Value::Array(["24h", "7d", "30d"].iter().map(|v| string(v)).collect()),
                ),
            ]),
        ),
        ("asOf", integer_schema()),
    ])
}
fn schema_for(name: &&str) -> Value {
    match *name {
        "expectedRevision" | "leaseGeneration" => integer_schema(),
        "repos" => array_schema(),
        _ => object(vec![("type", string("string"))]),
    }
}

fn integer_schema() -> Value {
    object(vec![("type", string("integer"))])
}
fn array_schema() -> Value {
    object(vec![
        ("type", string("array")),
        ("items", object(vec![("type", string("string"))])),
    ])
}
fn bounded_integer_schema(minimum: i64, maximum: i64) -> Value {
    object(vec![
        ("type", string("integer")),
        ("minimum", Value::Int(minimum)),
        ("maximum", Value::Int(maximum)),
    ])
}

fn state_array_schema() -> Value {
    object(vec![
        ("type", string("array")),
        (
            "items",
            object(vec![
                ("type", string("string")),
                (
                    "enum",
                    Value::Array(STATES.iter().map(|state| string(state)).collect()),
                ),
            ]),
        ),
    ])
}

fn outcome_patch(value: &Value) -> Result<OutcomePatch, String> {
    Ok(OutcomePatch {
        title: field(value, "title").map(str::to_string),
        description: field(value, "description").map(str::to_string),
        projects: value
            .get("projects")
            .map(|_| string_array(value, "projects"))
            .transpose()?,
        repos: value
            .get("repos")
            .map(|_| string_array(value, "repos"))
            .transpose()?,
        owner: field(value, "owner").map(str::to_string),
        target_metric: field(value, "targetMetric").map(str::to_string),
        baseline: value.get("baseline").cloned(),
        target: value.get("target").cloned(),
        observation_start: value.get("observationStart").and_then(Value::as_i64),
        observation_end: value.get("observationEnd").and_then(Value::as_i64),
        result_links: value
            .get("resultLinks")
            .map(|_| string_array(value, "resultLinks"))
            .transpose()?,
    })
}

fn required<'a>(value: &'a Value, key: &str) -> Result<&'a str, String> {
    field(value, key).ok_or_else(|| format!("{key} is required"))
}

fn integer(value: &Value, key: &str) -> Result<i64, String> {
    value
        .get(key)
        .and_then(Value::as_i64)
        .ok_or_else(|| format!("{key} must be an integer"))
}

fn lease_proof(value: &Value) -> Result<Option<LeaseProof<'_>>, String> {
    match (
        field(value, "leaseToken"),
        value.get("leaseGeneration").and_then(Value::as_i64),
    ) {
        (None, None) => Ok(None),
        (Some(token), Some(generation)) => Ok(Some(LeaseProof { token, generation })),
        _ => Err("leaseToken and leaseGeneration must be supplied together".into()),
    }
}

fn string_array(value: &Value, key: &str) -> Result<Vec<String>, String> {
    match value.get(key) {
        None => Ok(Vec::new()),
        Some(Value::Array(values)) => values
            .iter()
            .map(|value| {
                value
                    .as_str()
                    .map(str::to_string)
                    .ok_or_else(|| format!("{key} must contain strings"))
            })
            .collect(),
        Some(other) => {
            if other == &Value::Null {
                return Ok(Vec::new());
            }
            let Some(raw) = other.as_str() else {
                return Err(format!("{key} must be an array"));
            };
            let trimmed = raw.trim();
            if trimmed.is_empty() {
                return Ok(Vec::new());
            }
            if trimmed.starts_with('[') {
                return match parse(trimmed) {
                    Ok(Value::Array(values)) => values
                        .iter()
                        .map(|value| {
                            value
                                .as_str()
                                .map(str::to_string)
                                .ok_or_else(|| format!("{key} must contain strings"))
                        })
                        .collect(),
                    _ => Err(format!("{key} must be an array")),
                };
            }
            Ok(trimmed
                .split(',')
                .map(str::trim)
                .filter(|value| !value.is_empty())
                .map(str::to_string)
                .collect())
        }
    }
}

fn query_from(value: &Value) -> Result<Query, String> {
    Ok(Query {
        projects: string_array(value, "projects")?,
        repos: string_array(value, "repos")?,
        states: string_array(value, "states")?,
        priorities: string_array(value, "priorities")?,
        owners: string_array(value, "owners")?,
        types: string_array(value, "types")?,
        labels: string_array(value, "labels")?,
        parent: field(value, "parent").map(str::to_string),
        cycle: field(value, "cycle").map(str::to_string),
        text: field(value, "text").map(str::to_string),
        blocked: value.get("blocked").and_then(Value::as_bool),
        created_after: value.get("createdAfter").and_then(Value::as_i64),
        created_before: value.get("createdBefore").and_then(Value::as_i64),
        updated_after: value.get("updatedAfter").and_then(Value::as_i64),
        updated_before: value.get("updatedBefore").and_then(Value::as_i64),
        sort: field(value, "sort").unwrap_or("id").to_string(),
        descending: value
            .get("descending")
            .and_then(Value::as_bool)
            .unwrap_or(false),
        limit: value
            .get("limit")
            .and_then(Value::as_i64)
            .unwrap_or(100)
            .try_into()
            .map_err(|_| "limit must be positive")?,
        byte_limit: value
            .get("byteLimit")
            .and_then(Value::as_i64)
            .unwrap_or(1_048_576)
            .try_into()
            .map_err(|_| "byteLimit must be positive")?,
        cursor: field(value, "cursor").map(str::to_string),
    })
}

fn rpc_result(id: Value, result: Value) -> String {
    object(vec![
        ("jsonrpc", string("2.0")),
        ("id", id),
        ("result", result),
    ])
    .to_json()
}

fn rpc_error(id: Value, code: i64, message: &str) -> String {
    object(vec![
        ("jsonrpc", string("2.0")),
        ("id", id),
        (
            "error",
            object(vec![
                ("code", Value::Int(code)),
                ("message", string(message)),
            ]),
        ),
    ])
    .to_json()
}

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

    #[test]
    fn exposes_and_calls_multi_project_tools() {
        let path = std::env::temp_dir().join(format!("tasks-mcp-{}.db", std::process::id()));
        let _ = std::fs::remove_file(&path);
        let mcp = Mcp::new(Store::open(&path).unwrap(), "agent".into());
        let initialize = mcp.handle_line(r#"{"jsonrpc":"2.0","id":1,"method":"initialize"}"#);
        assert!(initialize.contains("2025-11-25"));
        let created = mcp.handle_line(r#"{"jsonrpc":"2.0","id":2,"method":"tools/call","params":{"name":"project_create","arguments":{"key":"OPS","name":"Ops","repos":["owner/repo"],"idempotencyKey":"p1"}}}"#);
        assert!(created.contains("OPS"));
        let listed = mcp.handle_line(r#"{"jsonrpc":"2.0","id":3,"method":"tools/call","params":{"name":"project_list","arguments":{}}}"#);
        assert!(listed.contains("owner/repo"));
        let tools = mcp.handle_line(r#"{"jsonrpc":"2.0","id":4,"method":"tools/list"}"#);
        for name in [
            "work_list",
            "work_block",
            "work_unblock",
            "events_read",
            "system_doctor",
            "system_backup",
        ] {
            assert!(tools.contains(&format!(r#""name":"{name}""#)));
        }
        let work = mcp.handle_line(r#"{"jsonrpc":"2.0","id":5,"method":"tools/call","params":{"name":"work_list","arguments":{"project":"OPS"}}}"#);
        assert!(work.contains(r#""content""#));
        let events = mcp.handle_line(r#"{"jsonrpc":"2.0","id":6,"method":"tools/call","params":{"name":"events_read","arguments":{"project":"OPS"}}}"#);
        assert!(events.contains(r#""content""#));
        let outcome = mcp.handle_line(r#"{"jsonrpc":"2.0","id":7,"method":"tools/call","params":{"name":"outcome_create","arguments":{"title":"Retention","projects":["OPS"],"owner":"owner","targetMetric":"retention","baseline":10,"target":20,"observationStart":1,"observationEnd":2,"idempotencyKey":"o1"}}}"#);
        assert!(outcome.contains("OUT-"));
        let metrics = mcp.handle_line(r#"{"jsonrpc":"2.0","id":8,"method":"tools/call","params":{"name":"metrics_query","arguments":{"projects":["OPS"],"period":"24h"}}}"#);
        assert!(metrics.contains("definitions"));
        let _ = std::fs::remove_file(path);
    }
    #[test]
    fn scoped_principal_cannot_read_other_projects_or_global_records() {
        let path = std::env::temp_dir().join(format!("tasks-mcp-scope-{}.db", std::process::id()));
        let _ = std::fs::remove_file(&path);
        let store = Store::open(&path).unwrap();
        store
            .bootstrap("admin", "Admin", "bootstrap", "bootstrap")
            .unwrap();
        store
            .create_project(
                "OPS",
                "Operations",
                vec!["owner/ops".into()],
                "admin",
                "ops",
            )
            .unwrap();
        store
            .create_project(
                "SECRET",
                "Secret",
                vec!["owner/secret".into()],
                "admin",
                "secret",
            )
            .unwrap();
        let secret = store
            .create_item(CreateItem {
                project: "SECRET",
                title: "Do not disclose",
                description: "",
                repo: Some("owner/secret"),
                priority: "normal",
                actor: "admin",
                idempotency_key: "secret-item",
            })
            .unwrap();
        store
            .create_item(CreateItem {
                project: "OPS",
                title: "Visible to ops",
                description: "",
                repo: Some("owner/ops"),
                priority: "normal",
                actor: "admin",
                idempotency_key: "ops-item",
            })
            .unwrap();
        store
            .upsert_principal(
                "ops-viewer",
                "agent",
                "Ops viewer",
                true,
                "admin",
                "principal",
            )
            .unwrap();
        store
            .grant_role("ops-viewer", "OPS", "viewer", "admin", "viewer-role")
            .unwrap();
        let mcp = Mcp::new(store, "ops-viewer".into());

        for (id, name, arguments) in [
            ("1", "project_list", "{}"),
            ("2", "work_query", "{}"),
            ("3", "events_read", "{}"),
            ("4", "system_doctor", "{}"),
        ] {
            let response = mcp.handle_line(&format!(
                r#"{{"jsonrpc":"2.0","id":{id},"method":"tools/call","params":{{"name":"{name}","arguments":{arguments}}}}}"#
            ));
            assert!(
                response.contains(r#""error""#),
                "{name} unexpectedly succeeded: {response}"
            );
        }
        let secret_id = field(&secret, "id").unwrap();
        for (id, name, arguments) in [
            ("5", "work_list", r#"{"project":"SECRET"}"#),
            ("6", "board_get", r#"{"project":"SECRET"}"#),
            ("7", "ready_frontier", r#"{"project":"SECRET"}"#),
            ("8", "events_read", r#"{"project":"SECRET"}"#),
            ("9", "work_get", r#"{"id":"SECRET_ID"}"#),
            ("10", "work_query", r#"{"projects":["SECRET"]}"#),
            (
                "11",
                "metrics_query",
                r#"{"projects":["SECRET"],"period":"24h"}"#,
            ),
        ] {
            let arguments = arguments.replace("SECRET_ID", secret_id);
            let response = mcp.handle_line(&format!(
                r#"{{"jsonrpc":"2.0","id":{id},"method":"tools/call","params":{{"name":"{name}","arguments":{arguments}}}}}"#
            ));
            assert!(
                response.contains(r#""error""#),
                "{name} unexpectedly succeeded: {response}"
            );
        }
        let allowed = mcp.handle_line(
            r#"{"jsonrpc":"2.0","id":12,"method":"tools/call","params":{"name":"work_query","arguments":{"projects":["OPS"]}}}"#,
        );
        assert!(
            allowed.contains(r#""result""#),
            "authorized project query failed: {allowed}"
        );
        let _ = std::fs::remove_file(path);
    }
    #[test]
    fn outcome_and_portfolio_schemas_match_supported_filters() {
        for name in ["outcome_query", "portfolio_query"] {
            let tool = tools()
                .into_iter()
                .find(|tool| field(tool, "name") == Some(name))
                .expect("outcome tool missing");
            let properties = tool
                .get("inputSchema")
                .and_then(|schema| schema.get("properties"))
                .expect("outcome properties missing");
            for supported in ["projects", "repos", "owners", "states", "limit"] {
                assert!(
                    properties.get(supported).is_some(),
                    "{name} lacks {supported}"
                );
            }
            for unsupported in [
                "priorities",
                "types",
                "labels",
                "cursor",
                "byteLimit",
                "text",
            ] {
                assert!(
                    properties.get(unsupported).is_none(),
                    "{name} advertises unsupported {unsupported}"
                );
            }
        }
    }
}