From e826b48a901d563f80a0f87024c89577f3c4a1a9 Mon Sep 17 00:00:00 2001 From: Rakesh Utekar <48244158+rakeshutekar@users.noreply.github.com> Date: Fri, 28 Aug 2026 14:47:07 -0700 Subject: [PATCH 1/3] fix(teams): enforce worker item visibility --- coworker/server/app.py | 6 +- coworker/server/manager.py | 12 +- coworker/teams/__init__.py | 2 + coworker/teams/dialect.py | 2 +- coworker/teams/model.py | 4 + coworker/teams/store.py | 103 +++++++++++---- tests/test_team_board.py | 60 ++++++++- tests/test_team_open_surface.py | 215 ++++++++++++++++++++++++++++++++ tests/test_team_store.py | 5 +- 9 files changed, 373 insertions(+), 36 deletions(-) diff --git a/coworker/server/app.py b/coworker/server/app.py index 2f5580853..f63599c15 100644 --- a/coworker/server/app.py +++ b/coworker/server/app.py @@ -164,6 +164,7 @@ def _connector_title(name: str) -> str: from .. import toolchain from ..teams.model import AuthorityError as TeamsAuthorityError from ..teams.model import BoardError as TeamsBoardError +from ..teams.model import BoardNotFoundError as TeamsBoardNotFoundError from .manager import SessionManager @@ -842,6 +843,8 @@ def _board(request: Request, handler): ) try: return handler(actor) + except TeamsBoardNotFoundError as error: + return JSONResponse({"error": str(error)}, status_code=404) except TeamsAuthorityError as error: return JSONResponse({"error": str(error)}, status_code=403) except (TeamsBoardError, ValueError) as error: @@ -873,7 +876,8 @@ def board_list_items( @app.get("/v1/board/item") def board_get_item(request: Request, space: str, id: int): return _board( - request, lambda actor: manager.team_store.get_item(space, int(id)) + request, + lambda actor: manager.team_store.get_item(space, int(id), actor=actor), ) @app.post("/v1/board/items") diff --git a/coworker/server/manager.py b/coworker/server/manager.py index 1486d9064..402857685 100644 --- a/coworker/server/manager.py +++ b/coworker/server/manager.py @@ -1697,7 +1697,9 @@ def board_item_detail(self, session_id: str, item_id: int) -> dict[str, Any]: if space is None: return {"error": "no board for this session"} try: - item = self.team_store.get_item(space, int(item_id)) + item = self.team_store.get_item( + space, int(item_id), actor=self._user_actor() + ) except TeamsBoardError as error: return {"error": str(error)} timeline: list[dict[str, Any]] = [] @@ -2295,7 +2297,9 @@ async def _drain_team_member( # item's ASSIGNEE — a filer merely hears about it. def _holds(event) -> bool: try: - item = self.team_store.get_item(team.space, int(event["item_id"])) + item = self.team_store.get_item( + team.space, int(event["item_id"]), actor=self._user_actor() + ) except Exception: return False return item["assignee"] == actor @@ -2388,7 +2392,9 @@ def _team_digest( item = None if item_id is not None: try: - item = self.team_store.get_item(team.space, int(item_id)) + item = self.team_store.get_item( + team.space, int(item_id), actor=self._user_actor() + ) except Exception: item = None title = f"#{item_id} {item['title']}" if item else f"#{item_id}" diff --git a/coworker/teams/__init__.py b/coworker/teams/__init__.py index 96320fd60..fd0806a4d 100644 --- a/coworker/teams/__init__.py +++ b/coworker/teams/__init__.py @@ -7,6 +7,7 @@ Actor, AuthorityError, BoardError, + BoardNotFoundError, ChainError, ItemState, Role, @@ -18,6 +19,7 @@ "Actor", "AuthorityError", "BoardError", + "BoardNotFoundError", "ChainError", "ItemState", "JournalStore", diff --git a/coworker/teams/dialect.py b/coworker/teams/dialect.py index 1bdb4d784..b96a01d81 100644 --- a/coworker/teams/dialect.py +++ b/coworker/teams/dialect.py @@ -149,7 +149,7 @@ def list_items( return self.store.list_items(space, self.actor, state=state, assignee=assignee) def get_item(self, space: str, item_id: int) -> dict[str, Any]: - return self.store.get_item(space, item_id) + return self.store.get_item(space, item_id, actor=self.actor) def create_item( self, diff --git a/coworker/teams/model.py b/coworker/teams/model.py index 2228e6613..3f89a34b8 100644 --- a/coworker/teams/model.py +++ b/coworker/teams/model.py @@ -85,6 +85,10 @@ class BoardError(Exception): """A verb call the board refuses — illegal transition, missing item, bad input.""" +class BoardNotFoundError(BoardError): + """A requested board object is missing or is not visible to the actor.""" + + class AuthorityError(BoardError): """The actor's role does not permit this verb on this item.""" diff --git a/coworker/teams/store.py b/coworker/teams/store.py index 876347bfb..547a19332 100644 --- a/coworker/teams/store.py +++ b/coworker/teams/store.py @@ -38,6 +38,7 @@ Actor, AuthorityError, BoardError, + BoardNotFoundError, ChainError, ItemState, Role, @@ -469,7 +470,16 @@ def create_item( ) with self._lock: if parent is not None: - parent_item = self._item(space, parent) + try: + parent_item = self._item(space, parent) + except BoardError: + raise BoardNotFoundError( + f"no visible item #{parent} in space {space!r}" + ) from None + if not self._item_visible_to(space, actor, parent_item): + raise BoardNotFoundError( + f"no visible item #{parent} in space {space!r}" + ) if case is None: case = parent_item["case_id"] or None item_id = self._next_item_id(space) @@ -489,7 +499,7 @@ def create_item( ) if self.journal is not None and case: self.journal.ensure_case(case, actor.id) - return self.get_item(space, item_id, seq=event["seq"]) + return self.get_item(space, item_id, actor=actor, seq=event["seq"]) def list_items( self, @@ -499,8 +509,11 @@ def list_items( state: Optional[str] = None, assignee: Optional[str] = None, ) -> list[dict[str, Any]]: - """Items in a space. Workers see only their slice: assigned items plus items - directly linked to those.""" + """Return actor-visible items in a space. + + A worker's slice contains items it owns or created, their direct links, + and—while claims are open—the open, unassigned claim pool. + """ where = ["space = ?"] params: list[Any] = [space] if state: @@ -517,32 +530,50 @@ def list_items( params, ).fetchall() items = [_row_to_item(row) for row in rows] + worker_slice = None + claims_open = None if actor.role == Role.WORKER: - visible = self._worker_slice(space, actor.id) - # On an open-claims board the claimable pool is visible too — a - # pull queue nobody can see is not a queue (drill-caught: an - # external worker with no assignment saw an empty board). Under - # lead-only policy workers can't act on it, so it stays hidden. + worker_slice = self._worker_slice(space, actor.id) claims_open = self.policy(space)["claims"] == "open" - items = [ - item - for item in items - if item["id"] in visible - or ( - claims_open - and item["state"] == ItemState.OPEN.value - and not item["assignee"] - ) - ] + items = [ + item + for item in items + if self._item_visible_to( + space, + actor, + item, + worker_slice=worker_slice, + claims_open=claims_open, + ) + ] for item in items: item["links"] = self._links_of(space, item["id"]) return items def get_item( - self, space: str, item_id: int, *, seq: Optional[int] = None + self, + space: str, + item_id: int, + *, + actor: Actor, + seq: Optional[int] = None, ) -> dict[str, Any]: + """Return one actor-visible item. + + The actor is required because detail reads enforce the same worker scope as + list reads. Missing and hidden items deliberately share one error contract. + """ with self._lock: - item = self._item(space, item_id) + try: + item = self._item(space, item_id) + except BoardError: + raise BoardNotFoundError( + f"no visible item #{item_id} in space {space!r}" + ) from None + if not self._item_visible_to(space, actor, item): + raise BoardNotFoundError( + f"no visible item #{item_id} in space {space!r}" + ) item["links"] = self._links_of(space, item_id) item["comments"] = self.comments(space, item_id) if seq is not None: @@ -587,7 +618,7 @@ def transition( }, taint=taint, ) - return self.get_item(space, item_id, seq=event["seq"]) + return self.get_item(space, item_id, actor=actor, seq=event["seq"]) def comment( self, @@ -652,7 +683,7 @@ def assign( assignee=assignee, previous=item["assignee"] or "", ) - return self.get_item(space, item_id, seq=event["seq"]) + return self.get_item(space, item_id, actor=actor, seq=event["seq"]) def claim(self, space: str, actor: Actor, item_id: int) -> dict[str, Any]: """Self-assign an open, unassigned item. Nobody stamps a claim — the store @@ -694,7 +725,7 @@ def claim(self, space: str, actor: Actor, item_id: int) -> dict[str, Any]: assignee=actor.id, previous="", ) - return self.get_item(space, item_id, seq=event["seq"]) + return self.get_item(space, item_id, actor=actor, seq=event["seq"]) def policy(self, space: str) -> dict[str, Any]: with self._lock: @@ -897,6 +928,30 @@ def _worker_slice(self, space: str, worker_id: str) -> set[int]: out.add(row["src"]) return out + def _item_visible_to( + self, + space: str, + actor: Actor, + item: dict[str, Any], + *, + worker_slice: Optional[set[int]] = None, + claims_open: Optional[bool] = None, + ) -> bool: + """Whether an actor may read one item through any board surface.""" + if actor.role != Role.WORKER: + return True + if worker_slice is None: + worker_slice = self._worker_slice(space, actor.id) + if item["id"] in worker_slice: + return True + if claims_open is None: + claims_open = self.policy(space)["claims"] == "open" + return ( + claims_open + and item["state"] == ItemState.OPEN.value + and not item["assignee"] + ) + def _links_of(self, space: str, item_id: int) -> list[dict[str, Any]]: rows = self._conn.execute( "SELECT src, kind, dst FROM team_links WHERE space = ?" diff --git a/tests/test_team_board.py b/tests/test_team_board.py index 42ce7e2df..66cedcf4d 100644 --- a/tests/test_team_board.py +++ b/tests/test_team_board.py @@ -3,6 +3,7 @@ import pytest from coworker.teams import Actor, AuthorityError, BoardError, Role, TeamStore +from coworker.teams.dialect import LocalDialect from coworker.teams.tools import board_tools USER = Actor(id="user", role=Role.USER) @@ -67,6 +68,23 @@ def test_child_inherits_parent_case(store): assert {"kind": "parent", "item": parent["id"]} in child["links"] +def test_worker_cannot_expand_its_slice_through_a_hidden_parent(store): + hidden = assigned_item(store, assignee="worker-2") + before = store.event_count(SPACE) + + with pytest.raises(BoardError): + store.create_item( + SPACE, + WORKER, + title="Bridge into another worker's task", + criteria="must not be created", + parent=hidden, + ) + + assert store.event_count(SPACE) == before + assert store.list_items(SPACE, WORKER) == [] + + # ------------------------------------------------------------------- transitions def test_full_happy_path(store): @@ -127,7 +145,7 @@ def test_rework_loop(store): store.transition(SPACE, WORKER, item_id, "in_progress") store.transition(SPACE, WORKER, item_id, "review") store.transition(SPACE, LEAD, item_id, "in_progress", comment="criteria 2 unmet") - item = store.get_item(SPACE, item_id) + item = store.get_item(SPACE, item_id, actor=LEAD) assert item["state"] == "in_progress" assert item["comments"][-1]["body"] == "criteria 2 unmet" @@ -161,9 +179,9 @@ def test_blocked_by_shows_on_the_other_side(store): a = store.create_item(SPACE, LEAD, title="A", criteria="c") b = store.create_item(SPACE, LEAD, title="B", criteria="c") store.link(SPACE, LEAD, a["id"], "blocks", b["id"]) - assert {"kind": "blocked_by", "item": a["id"]} in store.get_item(SPACE, b["id"])[ - "links" - ] + assert {"kind": "blocked_by", "item": a["id"]} in store.get_item( + SPACE, b["id"], actor=LEAD + )["links"] def test_artifact_refs_accumulate_on_the_item(store): @@ -174,10 +192,10 @@ def test_artifact_refs_accumulate_on_the_item(store): SPACE, WORKER, item_id, "review", comment="done", refs=["branch:fix/acl", "report:posture.html"], ) - item = store.get_item(SPACE, item_id) + item = store.get_item(SPACE, item_id, actor=WORKER) assert item["refs"] == ["branch:fix/acl", "report:posture.html"] # deduped, ordered store.rebuild(SPACE) - assert store.get_item(SPACE, item_id)["refs"] == item["refs"] + assert store.get_item(SPACE, item_id, actor=WORKER)["refs"] == item["refs"] # ------------------------------------------------------------- worker visibility @@ -194,6 +212,36 @@ def test_worker_sees_only_its_slice(store): assert len(store.list_items(SPACE, USER)) == 3 +def test_worker_item_detail_matches_list_visibility(store): + mine = assigned_item(store, assignee="worker-1") + theirs = assigned_item(store, assignee="worker-2") + claimable = store.create_item(SPACE, LEAD, title="Available", criteria="c") + linked = store.create_item(SPACE, LEAD, title="Dependency", criteria="c") + store.link(SPACE, LEAD, linked["id"], "blocks", mine) + worker = LocalDialect(store, journal=None, actor=WORKER) + + assert worker.get_item(SPACE, mine)["id"] == mine + assert worker.get_item(SPACE, linked["id"])["id"] == linked["id"] + assert worker.get_item(SPACE, claimable["id"])["id"] == claimable["id"] + assert store.get_item(SPACE, mine, actor=WORKER)["id"] == mine + + with pytest.raises(BoardError): + worker.get_item(SPACE, theirs) + with pytest.raises(BoardError): + store.get_item(SPACE, theirs, actor=WORKER) + + store.claim(SPACE, OTHER, claimable["id"]) + with pytest.raises(BoardError): + worker.get_item(SPACE, claimable["id"]) + + held = store.create_item(SPACE, LEAD, title="Held", criteria="c") + store.set_policy(SPACE, LEAD, claims="lead-only") + with pytest.raises(BoardError): + worker.get_item(SPACE, held["id"]) + + assert worker.get_item(SPACE, mine)["id"] == mine + + def test_worker_comments_only_on_its_slice(store): theirs = assigned_item(store, assignee="worker-2") with pytest.raises(AuthorityError, match="assigned"): diff --git a/tests/test_team_open_surface.py b/tests/test_team_open_surface.py index ad7b5d402..2ed8994d6 100644 --- a/tests/test_team_open_surface.py +++ b/tests/test_team_open_surface.py @@ -248,6 +248,166 @@ def test_token_binds_identity_and_store_enforces_authority(api): assert bad.status_code == 403 or bad.status_code == 400 +def test_board_item_reads_hide_foreign_worker_items(api): + client, manager, _ = api + lead_token = _tokens(manager).mint("lead-1", "lead") + nia_token = _tokens(manager).mint("nia", "worker") + webb_token = _tokens(manager).mint("webb", "worker") + lead = {"Authorization": f"Bearer {lead_token}"} + nia = {"Authorization": f"Bearer {nia_token}"} + webb = {"Authorization": f"Bearer {webb_token}"} + + def create(title, *, description=""): + response = client.post( + "/v1/board/items", + headers=lead, + json={ + "space": "proj", + "title": title, + "description": description, + "criteria": "done", + }, + ) + assert response.status_code == 200 + return response.json() + + created = create( + "Webb private investigation", description="credential rotation details" + ) + assigned_response = client.post( + "/v1/board/items/assign", + headers=lead, + json={"space": "proj", "id": created["id"], "assignee": "webb"}, + ) + assert assigned_response.status_code == 200 + + mine = create("Nia task") + assert client.post( + "/v1/board/items/assign", + headers=lead, + json={"space": "proj", "id": mine["id"], "assignee": "nia"}, + ).status_code == 200 + claimable = create("Available") + linked = create("Linked dependency") + assert client.post( + "/v1/board/link", + headers=lead, + json={ + "space": "proj", + "src": linked["id"], + "kind": "blocks", + "dst": mine["id"], + }, + ).status_code == 200 + + denied = client.get( + "/v1/board/item", + headers=nia, + params={"space": "proj", "id": created["id"]}, + ) + assert denied.status_code == 404 + assert "Webb private investigation" not in denied.text + assert "credential rotation details" not in denied.text + + for item_id in (mine["id"], claimable["id"], linked["id"]): + assert client.get( + "/v1/board/item", + headers=nia, + params={"space": "proj", "id": item_id}, + ).status_code == 200 + + assert client.post( + "/v1/board/items/claim", + headers=webb, + json={"space": "proj", "id": claimable["id"]}, + ).status_code == 200 + assert client.get( + "/v1/board/item", + headers=nia, + params={"space": "proj", "id": claimable["id"]}, + ).status_code == 404 + + held = create("Held") + assert client.post( + "/v1/board/policy", + headers=lead, + json={"space": "proj", "claims": "lead-only"}, + ).status_code == 200 + assert client.get( + "/v1/board/item", + headers=nia, + params={"space": "proj", "id": held["id"]}, + ).status_code == 404 + assert client.get( + "/v1/board/item", + headers=nia, + params={"space": "proj", "id": mine["id"]}, + ).status_code == 200 + + allowed = client.get( + "/v1/board/item", + headers=lead, + params={"space": "proj", "id": created["id"]}, + ) + assert allowed.status_code == 200 + assert allowed.json()["title"] == "Webb private investigation" + + +def test_board_item_reads_return_not_found_for_missing_items(api): + client, manager, _ = api + nia_token = _tokens(manager).mint("nia", "worker") + lead_token = _tokens(manager).mint("lead-1", "lead") + user_token = _tokens(manager).mint("user", "user") + nia = {"Authorization": f"Bearer {nia_token}"} + lead = {"Authorization": f"Bearer {lead_token}"} + user = {"Authorization": f"Bearer {user_token}"} + + for headers in (nia, lead, user): + missing = client.get( + "/v1/board/item", + headers=headers, + params={"space": "proj", "id": 999}, + ) + assert missing.status_code == 404 + + +def test_board_item_create_hides_missing_and_foreign_parents(api): + client, manager, _ = api + lead_token = _tokens(manager).mint("lead-1", "lead") + nia_token = _tokens(manager).mint("nia", "worker") + lead = {"Authorization": f"Bearer {lead_token}"} + nia = {"Authorization": f"Bearer {nia_token}"} + + hidden_response = client.post( + "/v1/board/items", + headers=lead, + json={"space": "proj", "title": "Webb task", "criteria": "done"}, + ) + assert hidden_response.status_code == 200 + hidden = hidden_response.json() + assert client.post( + "/v1/board/items/assign", + headers=lead, + json={"space": "proj", "id": hidden["id"], "assignee": "webb"}, + ).status_code == 200 + before = manager.team_store.event_count("proj") + + for parent in (hidden["id"], 999): + denied = client.post( + "/v1/board/items", + headers=nia, + json={ + "space": "proj", + "title": "Probe", + "criteria": "must not be created", + "parent": parent, + }, + ) + assert denied.status_code == 404 + + assert manager.team_store.event_count("proj") == before + + def test_remote_dialect_round_trip(api): client, manager, app = api lead_token = _tokens(manager).mint("lead-1", "lead") @@ -288,6 +448,13 @@ def test_remote_dialect_round_trip(api): with pytest.raises(BoardError, match="only open items"): nia.claim("proj", item["id"]) + foreign = lead.create_item( + "proj", title="Webb-only task", criteria="visible only to webb" + ) + lead.assign("proj", foreign["id"], "webb") + with pytest.raises(BoardError, match="no visible item"): + nia.get_item("proj", foreign["id"]) + # policy flip over the wire blocks the next worker claim second = lead.create_item("proj", title="Held back", criteria="c") lead.set_policy("proj", claims="lead-only") @@ -477,10 +644,33 @@ def call(name, arguments): call("board_claim", {"item": item["id"]}) call("board_move", {"item": item["id"], "to": "in_progress"}) + payload = json.loads(call("board_show", {"item": item["id"]})[0].text) + assert (payload["id"], payload["assignee"]) == (item["id"], "nia") shown = lead_dialect.get_item("proj", item["id"]) assert (shown["assignee"], shown["state"]) == ("nia", "in_progress") +def test_mcp_board_show_does_not_return_a_foreign_item(tmp_path): + import anyio + + from coworker.teams.mcp_server import build + + lead = local_dialect(tmp_path, actor="lead-1", role="lead") + foreign = lead.create_item( + "proj", title="Private Webb task", criteria="not visible to nia" + ) + lead.assign("proj", foreign["id"], "webb") + worker = build(LocalDialect(lead.store, lead.journal, NIA), space="proj") + + result = anyio.run( + lambda: worker.call_tool("board_show", {"item": foreign["id"]}) + ) + payload = json.loads(result[0].text) + + assert "error" in payload + assert "Private Webb task" not in payload["error"] + + # ------------------------------------------------------------------ CLI @@ -497,6 +687,10 @@ def test_cli_headless_flow(tmp_path, capsys): ["board", "claim", "1", *space_args, "--actor", "nia", "--role", "worker"] ) == 0 assert "claimed #1" in capsys.readouterr().out + assert main( + ["board", "show", "1", *space_args, "--actor", "nia", "--role", "worker"] + ) == 0 + assert "CLI item" in capsys.readouterr().out assert main(["board", "list", *space_args, "--json"]) == 0 items = json.loads(capsys.readouterr().out) assert [(i["id"], i["assignee"]) for i in items] == [(1, "nia")] @@ -518,6 +712,27 @@ def test_cli_headless_flow(tmp_path, capsys): assert "found it" in capsys.readouterr().out +def test_cli_worker_cannot_show_a_foreign_item(tmp_path, capsys): + from coworker.teams.cli import main + + space_args = ["--db", str(tmp_path), "--space", "proj"] + lead_args = [*space_args, "--actor", "lead-1", "--role", "lead"] + worker_args = [*space_args, "--actor", "nia", "--role", "worker"] + + assert main( + ["board", "create", "Private Webb task", "--criteria", "c", *lead_args] + ) == 0 + capsys.readouterr() + assert main(["board", "assign", "1", "webb", *lead_args]) == 0 + capsys.readouterr() + + assert main(["board", "show", "1", *worker_args]) == 1 + output = capsys.readouterr() + assert output.out == "" + assert "no visible item" in output.err + assert "Private Webb task" not in output.err + + def test_cli_token_mint_and_list(tmp_path, capsys): from coworker.teams.cli import main diff --git a/tests/test_team_store.py b/tests/test_team_store.py index 2269ca88e..b7e04bcc6 100644 --- a/tests/test_team_store.py +++ b/tests/test_team_store.py @@ -95,7 +95,10 @@ def test_rebuild_reproduces_the_projection(store): assert after[0]["state"] == "in_progress" assert after[0]["assignee"] == "worker-1" # created_ts comes from the event, so replay is deterministic - assert store.get_item("proj", item["id"])["created_ts"] == item["created_ts"] + assert ( + store.get_item("proj", item["id"], actor=USER)["created_ts"] + == item["created_ts"] + ) def test_taint_travels_with_the_record(store): From 0939fd6624a72c2adda678aa4e18d3f96ce14853 Mon Sep 17 00:00:00 2001 From: Rakesh Utekar <48244158+rakeshutekar@users.noreply.github.com> Date: Fri, 28 Aug 2026 15:35:55 -0700 Subject: [PATCH 2/3] fix(teams): authorize attachment reads Require attachment reads to resolve through an actor-visible item in the requested board space. Record authoritative attachment provenance so forged comment or transition refs cannot grant blob access, while preserving legacy refs through an atomic one-time migration. BREAKING CHANGE: BoardDialect.attachment and the /v1/board/attachment endpoint now require a board space. --- coworker/server/app.py | 13 +- coworker/server/manager.py | 11 + coworker/teams/attachments.py | 16 +- coworker/teams/cli.py | 2 +- coworker/teams/dialect.py | 19 +- coworker/teams/store.py | 175 ++++++++++++++- coworker/teams/tools.py | 4 +- tests/test_team_open_surface.py | 370 +++++++++++++++++++++++++++++++- tests/test_team_wake.py | 85 ++++++++ 9 files changed, 667 insertions(+), 28 deletions(-) diff --git a/coworker/server/app.py b/coworker/server/app.py index f63599c15..457c08110 100644 --- a/coworker/server/app.py +++ b/coworker/server/app.py @@ -783,12 +783,12 @@ def session_board_attachment(session_id: str, name: str): from fastapi.responses import Response try: - path = manager.attachment_store.path_for(name) + data, mime = manager.board_attachment(session_id, name) except TeamsBoardError as error: return JSONResponse({"error": str(error)}, status_code=404) return Response( - content=path.read_bytes(), - media_type=manager.attachment_store.mime_for(name), + content=data, + media_type=mime, ) @app.post("/v1/sessions/{session_id}/board/comment") @@ -998,22 +998,23 @@ def run(actor): data, str(body.get("filename", "")) ) filename = str(body.get("filename", "")) - event = manager.team_store.comment( + event = manager.team_store.attach_ref( str(body.get("space", "")), actor, int(body.get("id", 0)), str(body.get("caption", "")) or f"attached {filename}", - refs=[ref], + ref, ) return {"ref": ref, "seq": event["seq"]} return _board(request, run) @app.get("/v1/board/attachment") - def board_attachment(request: Request, name: str): + def board_attachment(request: Request, name: str, space: str): def run(actor): from fastapi.responses import Response + manager.team_store.require_attachment_access(space, actor, name) path = manager.attachment_store.path_for(name) return Response( content=path.read_bytes(), diff --git a/coworker/server/manager.py b/coworker/server/manager.py index 402857685..130ede1f2 100644 --- a/coworker/server/manager.py +++ b/coworker/server/manager.py @@ -1688,6 +1688,17 @@ def _board_space(self, session_id: str) -> Optional[str]: def _user_actor(self) -> TeamActor: return TeamActor(id="user", role=TeamRole.USER) + def board_attachment(self, session_id: str, stored: str) -> tuple[bytes, str]: + """Read an attachment referenced by the session's board as the user.""" + space = self._board_space(session_id) + if space is None: + raise TeamsBoardError("attachment not found") + self.team_store.require_attachment_access( + space, self._user_actor(), stored + ) + path = self.attachment_store.path_for(stored) + return path.read_bytes(), self.attachment_store.mime_for(stored) + def board_item_detail(self, session_id: str, item_id: int) -> dict[str, Any]: """One item in full, with its TIMELINE — creations, assignments, transitions, and comments merged chronologically (the detail pane renders diff --git a/coworker/teams/attachments.py b/coworker/teams/attachments.py index 807981c8b..90acb9fe1 100644 --- a/coworker/teams/attachments.py +++ b/coworker/teams/attachments.py @@ -22,7 +22,7 @@ from pathlib import Path from typing import Optional -from .model import BoardError +from .model import BoardError, BoardNotFoundError ATTACHMENT_SCHEME = "attachment://" MAX_ATTACHMENT_BYTES = 10 * 1024 * 1024 @@ -69,12 +69,10 @@ def put(self, data: bytes, filename: str) -> str: def path_for(self, stored: str) -> Path: """Resolve a stored name (`.`) to its file. The strict name check is the traversal guard — nothing else reaches the filesystem.""" - stored = stored.strip() - if not _STORED_NAME.fullmatch(stored): - raise BoardError(f"not an attachment name: {stored!r}") + stored = validate_stored_name(stored) path = self.root / stored if not path.exists(): - raise BoardError(f"no attachment {stored}") + raise BoardNotFoundError("attachment not found") return path def mime_for(self, stored: str) -> str: @@ -88,6 +86,14 @@ def stored_name(ref: str) -> Optional[str]: return ref[len(ATTACHMENT_SCHEME):].split("#", 1)[0] +def validate_stored_name(stored: str) -> str: + """Return one normalized stored name, rejecting malformed input.""" + stored = stored.strip() + if not _STORED_NAME.fullmatch(stored): + raise BoardError(f"not an attachment name: {stored!r}") + return stored + + def _validate(data: bytes, filename: str) -> str: if not data: raise BoardError("attachment is empty") diff --git a/coworker/teams/cli.py b/coworker/teams/cli.py index 6bb5de266..98d730825 100644 --- a/coworker/teams/cli.py +++ b/coworker/teams/cli.py @@ -356,7 +356,7 @@ def _cmd_attachment(args) -> int: from .attachments import stored_name stored = stored_name(args.ref) or args.ref - data, _mime = _dialect(args).attachment(stored) + data, _mime = _dialect(args).attachment(_space(args), stored) out = Path(args.out) if args.out else Path( args.ref.rsplit("#", 1)[-1] if "#" in args.ref else stored ) diff --git a/coworker/teams/dialect.py b/coworker/teams/dialect.py index b96a01d81..df66ac6e6 100644 --- a/coworker/teams/dialect.py +++ b/coworker/teams/dialect.py @@ -87,7 +87,9 @@ def attach( *, caption: str = "", ) -> dict[str, Any]: ... - def attachment(self, stored: str) -> tuple[bytes, str]: ... + def attachment(self, space: str, stored: str) -> tuple[bytes, str]: + """Read a blob referenced by an actor-visible item in ``space``.""" + ... def policy(self, space: str) -> dict[str, Any]: ... def set_policy(self, space: str, *, claims: str) -> dict[str, Any]: ... def pending(self, space: str, *, limit: int = 200) -> list[dict[str, Any]]: ... @@ -212,22 +214,23 @@ def attach( *, caption: str = "", ) -> dict[str, Any]: - # Attach = store blob + a normal comment event carrying the ref. Comment + # Attach = store blob + an attributed attachment-comment event. Comment # authority IS attach authority (workers attach on their slice only). if self.attachments is None: raise BoardError("no attachment store is attached to this board") ref = self.attachments.put(data, filename) - return self.store.comment( + return self.store.attach_ref( space, self.actor, item_id, caption or f"attached {filename}", - refs=[ref], + ref, ) - def attachment(self, stored: str) -> tuple[bytes, str]: + def attachment(self, space: str, stored: str) -> tuple[bytes, str]: if self.attachments is None: raise BoardError("no attachment store is attached to this board") + self.store.require_attachment_access(space, self.actor, stored) path = self.attachments.path_for(stored) return path.read_bytes(), self.attachments.mime_for(stored) @@ -458,8 +461,10 @@ def attach( }, ) - def attachment(self, stored: str) -> tuple[bytes, str]: - response = self._client.get("/v1/board/attachment", params={"name": stored}) + def attachment(self, space: str, stored: str) -> tuple[bytes, str]: + response = self._client.get( + "/v1/board/attachment", params={"space": space, "name": stored} + ) if response.status_code >= 400: self._unwrap(response) # raises with the server's message return response.content, response.headers.get( diff --git a/coworker/teams/store.py b/coworker/teams/store.py index 547a19332..188a9ee01 100644 --- a/coworker/teams/store.py +++ b/coworker/teams/store.py @@ -31,6 +31,7 @@ from pathlib import Path from typing import Any, Optional +from .attachments import stored_name, validate_stored_name from .model import ( EDGES, LINK_KINDS, @@ -51,6 +52,7 @@ # external. "lead-only" turns claims off; assignment stays with the lead/user. A lead # on an open board can still reserve individual items by assigning them to itself. CLAIM_POLICIES = ("open", "lead-only") +ATTACHMENT_REFS_MIGRATION = "attachment_refs_v1" # Event kinds. Chat lands later with the chat surface; the record shape already fits. # Journal entries live in their own case-keyed store (teams.journal) — cases outlive @@ -137,6 +139,13 @@ def __init__(self, db_path: str | Path, *, journal: Any = None) -> None: dst INTEGER NOT NULL, UNIQUE (space, src, kind, dst) ); + CREATE TABLE IF NOT EXISTS team_attachment_refs ( + space TEXT NOT NULL, + stored TEXT NOT NULL, + item_id INTEGER NOT NULL, + event_seq INTEGER NOT NULL, + PRIMARY KEY (space, stored, item_id) + ); CREATE TABLE IF NOT EXISTS team_meta ( space TEXT PRIMARY KEY, head_hash TEXT NOT NULL, @@ -150,8 +159,25 @@ def __init__(self, db_path: str | Path, *, journal: Any = None) -> None: space TEXT PRIMARY KEY, claims TEXT NOT NULL DEFAULT 'open' ); + CREATE TABLE IF NOT EXISTS team_migrations ( + name TEXT PRIMARY KEY + ); """) - self._conn.commit() + migrated = self._conn.execute( + "SELECT 1 FROM team_migrations WHERE name = ?", + (ATTACHMENT_REFS_MIGRATION,), + ).fetchone() + if migrated is None: + try: + self._backfill_legacy_attachment_refs() + self._conn.execute( + "INSERT INTO team_migrations (name) VALUES (?)", + (ATTACHMENT_REFS_MIGRATION,), + ) + self._conn.commit() + except Exception: + self._conn.rollback() + raise # ------------------------------------------------------------------ events core @@ -423,6 +449,11 @@ def rebuild(self, space: str) -> None: with self._lock: self._conn.execute("DELETE FROM team_items WHERE space = ?", (space,)) self._conn.execute("DELETE FROM team_links WHERE space = ?", (space,)) + self._conn.execute( + "DELETE FROM team_attachment_refs" + " WHERE space = ? AND event_seq != 0", + (space,), + ) rows = self._conn.execute( "SELECT seq, ts, kind, actor, item_id, payload FROM team_events" " WHERE space = ? ORDER BY seq", @@ -580,6 +611,84 @@ def get_item( item["seq"] = seq return item + def require_attachment_access( + self, space: str, actor: Actor, stored: str + ) -> None: + """Require an actor-visible item to carry an authoritative attachment.""" + stored = validate_stored_name(stored) + with self._lock: + rows = self._conn.execute( + "SELECT item.* FROM team_items AS item" + " JOIN team_attachment_refs AS attachment" + " ON attachment.space = item.space" + " AND attachment.item_id = item.id" + " WHERE attachment.space = ? AND attachment.stored = ?" + " ORDER BY item.id", + (space, stored), + ).fetchall() + worker_slice = None + claims_open = None + if actor.role == Role.WORKER: + worker_slice = self._worker_slice(space, actor.id) + claims_open = self.policy(space)["claims"] == "open" + for row in rows: + item = _row_to_item(row) + if not self._item_visible_to( + space, + actor, + item, + worker_slice=worker_slice, + claims_open=claims_open, + ): + continue + return + raise BoardNotFoundError("attachment not found") + + def attach_ref( + self, + space: str, + actor: Actor, + item_id: int, + body: str, + ref: str, + *, + taint: bool = False, + ) -> dict[str, Any]: + """Attach one stored blob through an attributed comment event. + + The dedicated payload field is the authoritative provenance marker; + arbitrary artifact refs on normal comments and transitions never grant + attachment-byte access. + """ + if not (body or "").strip(): + raise BoardError("comment body is required") + stored = stored_name(ref) + if stored is None: + raise BoardError(f"not an attachment ref: {ref!r}") + stored = validate_stored_name(stored) + with self._lock: + item = self._item(space, item_id) + if actor.role == Role.WORKER and item_id not in self._worker_slice( + space, actor.id + ): + raise AuthorityError( + f"worker {actor.id} may only comment on its assigned items" + " and items linked to them" + ) + return self.append_event( + space, + ITEM_COMMENTED, + actor, + item_id=item_id, + case_id=item["case_id"] or None, + payload={ + "body": body, + "refs": [ref], + "attachments": [stored], + }, + taint=taint, + ) + def transition( self, space: str, @@ -812,9 +921,9 @@ def _apply( item_id: Optional[int], payload: dict[str, Any], ) -> None: - """Fold one event into the projections. The ONLY writer of team_items and - team_links — shared by live appends and rebuild(), so replay always - reproduces the materialized state.""" + """Fold one event into the projections. The ONLY writer of item, link, + and attachment-reference projections — shared by live appends and + rebuild(), so replay always reproduces the materialized state.""" if kind == ITEM_CREATED: self._conn.execute( """ @@ -851,6 +960,9 @@ def _apply( self._merge_refs(space, item_id, payload.get("refs")) elif kind == ITEM_COMMENTED: self._merge_refs(space, item_id, payload.get("refs")) + self._merge_attachment_refs( + space, item_id, payload.get("attachments"), seq + ) elif kind == ITEM_ASSIGNED: self._conn.execute( "UPDATE team_items SET assignee = ?, updated_seq = ?" @@ -865,7 +977,8 @@ def _apply( ) # Comment bodies and journal entries have no materialized state: their # projections read straight off the (indexed) log. Only the artifact - # refs a comment carries fold onto the item. + # refs a comment carries fold onto the item; authoritative attachment + # markers also fold into their indexed projection. def _merge_refs( self, space: str, item_id: Optional[int], refs: Optional[list] @@ -885,6 +998,51 @@ def _merge_refs( (json.dumps(merged), space, item_id), ) + def _merge_attachment_refs( + self, + space: str, + item_id: Optional[int], + attachments: Optional[list], + seq: int, + ) -> None: + if not attachments or item_id is None: + return + for stored in attachments: + self._conn.execute( + "INSERT OR IGNORE INTO team_attachment_refs" + " (space, stored, item_id, event_seq) VALUES (?, ?, ?, ?)", + (space, validate_stored_name(str(stored)), item_id, seq), + ) + + def _backfill_legacy_attachment_refs(self) -> None: + """Snapshot pre-provenance refs once when the projection is introduced. + + Old attach events were indistinguishable from generic comment refs. Rows + grandfathered at upgrade use event_seq=0 so rebuild preserves that fixed + compatibility boundary; refs added after the migration are never inferred. + """ + rows = self._conn.execute( + "SELECT space, id, refs FROM team_items" + ).fetchall() + for row in rows: + try: + refs = json.loads(row["refs"] or "[]") + except (TypeError, json.JSONDecodeError): + continue + for ref in refs: + stored = stored_name(str(ref)) + if stored is None: + continue + try: + stored = validate_stored_name(stored) + except BoardError: + continue + self._conn.execute( + "INSERT OR IGNORE INTO team_attachment_refs" + " (space, stored, item_id, event_seq) VALUES (?, ?, ?, 0)", + (row["space"], stored, row["id"]), + ) + def _check_transition_authority( self, actor: Actor, item: dict[str, Any], current: ItemState, target: ItemState ) -> None: @@ -1047,7 +1205,12 @@ def rekey_space(self, old: str, new: str) -> bool: (new, prev, record["hash"], row["seq"]), ) prev = record["hash"] - for table in ("team_items", "team_links", "team_settings"): + for table in ( + "team_items", + "team_links", + "team_attachment_refs", + "team_settings", + ): self._conn.execute( f"UPDATE {table} SET space = ? WHERE space = ?", (new, old) ) diff --git a/coworker/teams/tools.py b/coworker/teams/tools.py index c704e2708..e004a037d 100644 --- a/coworker/teams/tools.py +++ b/coworker/teams/tools.py @@ -165,12 +165,12 @@ def attach_image(item: int, path: str, caption: str = "") -> dict: except (BoardError, ValueError) as error: return {"error": str(error)} return _call( - store.comment, + store.attach_ref, space, actor, item, caption or f"attached {source.name}", - refs=[ref], + ref, taint=taint(), ) diff --git a/tests/test_team_open_surface.py b/tests/test_team_open_surface.py index 2ed8994d6..6cde40d58 100644 --- a/tests/test_team_open_surface.py +++ b/tests/test_team_open_surface.py @@ -2,6 +2,7 @@ seam (local and remote), token-bound identity on `/v1/board`, and the MCP/CLI front doors.""" +import base64 import json import pytest @@ -551,8 +552,12 @@ def test_attach_over_the_wire_and_fetch(api): # and the lead can fetch the bytes back from coworker.teams.attachments import stored_name - data, mime = lead.attachment(stored_name(result["ref"])) + stored = stored_name(result["ref"]) + data, mime = lead.attachment("proj", stored) assert data == PNG and mime == "image/png" + assert nia.attachment("proj", stored) == (PNG, "image/png") + with pytest.raises(BoardError, match="attachment not found"): + lead.attachment("another-space", stored) # a worker cannot attach to an item outside its slice other = lead.create_item("proj", title="Not nia's", criteria="c") @@ -561,6 +566,300 @@ def test_attach_over_the_wire_and_fetch(api): nia.attach("proj", other["id"], PNG, "sneaky.png") +def test_board_attachment_read_hides_foreign_worker_reference(api): + client, manager, app = api + from fastapi.testclient import TestClient + from coworker.teams.attachments import stored_name + + lead_token = _tokens(manager).mint("lead-1", "lead") + nia_token = _tokens(manager).mint("nia", "worker") + lead = {"Authorization": f"Bearer {lead_token}"} + nia = {"Authorization": f"Bearer {nia_token}"} + + item_response = client.post( + "/v1/board/items", + headers=lead, + json={"space": "proj", "title": "Webb evidence", "criteria": "c"}, + ) + assert item_response.status_code == 200 + item = item_response.json() + assert client.post( + "/v1/board/items/assign", + headers=lead, + json={"space": "proj", "id": item["id"], "assignee": "webb"}, + ).status_code == 200 + attached = client.post( + "/v1/board/items/attach", + headers=lead, + json={ + "space": "proj", + "id": item["id"], + "filename": "private.png", + "data_b64": base64.b64encode(PNG).decode("ascii"), + }, + ) + assert attached.status_code == 200 + + stored = stored_name(attached.json()["ref"]) + denied = client.get( + "/v1/board/attachment", + headers=nia, + params={"space": "proj", "name": stored}, + ) + + assert denied.status_code == 404 + assert denied.json() == {"error": "attachment not found"} + + remote = RemoteDialect( + "http://board.test", + nia_token, + client=TestClient(app, base_url="http://board.test"), + ) + with pytest.raises(BoardError, match="attachment not found"): + remote.attachment("proj", stored) + remote.close() + + +def test_board_attachment_read_requires_token(api): + client, _, _ = api + + response = client.get( + "/v1/board/attachment", + params={"space": "proj", "name": f"{'0' * 64}.png"}, + ) + + assert response.status_code == 401 + assert "board token required" in response.json()["error"] + + +def test_board_attachment_read_rejects_malformed_name(api): + client, manager, _ = api + lead = { + "Authorization": f"Bearer {_tokens(manager).mint('lead-1', 'lead')}" + } + + response = client.get( + "/v1/board/attachment", + headers=lead, + params={"space": "proj", "name": "../../private.png"}, + ) + + assert response.status_code == 400 + assert "not an attachment name" in response.json()["error"] + + +def test_board_attachment_read_hides_unreferenced_blob(api): + client, manager, _ = api + from coworker.teams.attachments import stored_name + + lead = { + "Authorization": f"Bearer {_tokens(manager).mint('lead-1', 'lead')}" + } + ref = manager.attachment_store.put(PNG, "orphan.png") + + response = client.get( + "/v1/board/attachment", + headers=lead, + params={"space": "proj", "name": stored_name(ref)}, + ) + + assert response.status_code == 404 + assert response.json() == {"error": "attachment not found"} + + +def test_board_attachment_read_hides_missing_referenced_blob(api): + client, manager, _ = api + from coworker.teams.attachments import stored_name + + lead = { + "Authorization": f"Bearer {_tokens(manager).mint('lead-1', 'lead')}" + } + item = client.post( + "/v1/board/items", + headers=lead, + json={"space": "proj", "title": "Evidence", "criteria": "c"}, + ).json() + attached = client.post( + "/v1/board/items/attach", + headers=lead, + json={ + "space": "proj", + "id": item["id"], + "filename": "missing.png", + "data_b64": base64.b64encode(PNG).decode("ascii"), + }, + ).json() + stored = stored_name(attached["ref"]) + manager.attachment_store.path_for(stored).unlink() + + response = client.get( + "/v1/board/attachment", + headers=lead, + params={"space": "proj", "name": stored}, + ) + + assert response.status_code == 404 + assert response.json() == {"error": "attachment not found"} + + +def test_local_attachment_read_tracks_worker_visibility(tmp_path): + from coworker.teams.attachments import stored_name + + lead = local_dialect(tmp_path, actor="lead-1", role="lead") + worker = LocalDialect( + lead.store, lead.journal, NIA, attachments=lead.attachments + ) + claimable = lead.create_item("proj", title="Available", criteria="c") + result = lead.attach("proj", claimable["id"], PNG, "available.png") + stored = stored_name(result["payload"]["refs"][0]) + + assert worker.attachment("proj", stored) == (PNG, "image/png") + + lead.set_policy("proj", claims="lead-only") + with pytest.raises(BoardError, match="attachment not found"): + worker.attachment("proj", stored) + + assigned = lead.create_item("proj", title="Assigned", criteria="c") + lead.assign("proj", assigned["id"], "nia") + assigned_ref = lead.attach( + "proj", assigned["id"], PNG, "assigned.png" + )["payload"]["refs"][0] + assert worker.attachment("proj", stored_name(assigned_ref)) == ( + PNG, + "image/png", + ) + + with pytest.raises(BoardError, match="attachment not found"): + lead.attachment("another-space", stored) + + +def test_local_attachment_read_allows_a_visible_deduplicated_reference(tmp_path): + from coworker.teams.attachments import stored_name + + lead = local_dialect(tmp_path, actor="lead-1", role="lead") + worker = LocalDialect( + lead.store, lead.journal, NIA, attachments=lead.attachments + ) + foreign = lead.create_item("proj", title="Webb evidence", criteria="c") + lead.assign("proj", foreign["id"], "webb") + foreign_ref = lead.attach("proj", foreign["id"], PNG, "foreign.png")["payload"][ + "refs" + ][0] + with pytest.raises(BoardError, match="attachment not found"): + worker.attachment("proj", stored_name(foreign_ref)) + + mine = worker.create_item("proj", title="Nia evidence", criteria="c") + own_ref = worker.attach("proj", mine["id"], PNG, "mine.png")["payload"]["refs"][ + 0 + ] + + assert stored_name(own_ref) == stored_name(foreign_ref) + assert worker.attachment("proj", stored_name(own_ref)) == (PNG, "image/png") + + lead.store.rebuild("proj") + assert worker.attachment("proj", stored_name(own_ref)) == (PNG, "image/png") + + assert lead.store.rekey_space("proj", "moved") is True + assert worker.attachment("moved", stored_name(own_ref)) == (PNG, "image/png") + + +def test_local_attachment_read_rejects_forged_comment_and_transition_refs(tmp_path): + from coworker.teams.attachments import stored_name + + lead = local_dialect(tmp_path, actor="lead-1", role="lead") + worker = LocalDialect( + lead.store, lead.journal, NIA, attachments=lead.attachments + ) + foreign = lead.create_item("proj", title="Webb evidence", criteria="c") + lead.assign("proj", foreign["id"], "webb") + comment_ref = lead.attach( + "proj", foreign["id"], PNG, "foreign-comment.png" + )["payload"]["refs"][0] + transition_ref = lead.attach( + "proj", foreign["id"], PNG + b"transition", "foreign-transition.png" + )["payload"]["refs"][0] + + mine = lead.create_item("proj", title="Nia task", criteria="c") + lead.assign("proj", mine["id"], "nia") + worker.comment("proj", mine["id"], "found this", refs=[comment_ref]) + worker.transition( + "proj", + mine["id"], + "in_progress", + refs=[transition_ref], + ) + + with pytest.raises(BoardError, match="attachment not found"): + worker.attachment("proj", stored_name(comment_ref)) + with pytest.raises(BoardError, match="attachment not found"): + worker.attachment("proj", stored_name(transition_ref)) + + +def test_legacy_attachment_refs_are_grandfathered_across_rebuild(tmp_path): + import sqlite3 + + from coworker.teams.attachments import AttachmentStore, stored_name + + db_path = tmp_path / "teams.db" + attachments = AttachmentStore(tmp_path / "attachments") + store = TeamStore(db_path) + item = store.create_item("proj", LEAD, title="Legacy", criteria="c") + ref = attachments.put(PNG, "legacy.png") + # Before provenance markers, attach was an ordinary comment carrying a ref. + store.comment("proj", LEAD, item["id"], "attached legacy.png", refs=[ref]) + store.close() + with sqlite3.connect(db_path) as connection: + # Simulate an old DB—or an interrupted upgrade whose projection table + # exists but whose atomic migration marker was never committed. + connection.execute( + "DELETE FROM team_migrations WHERE name = 'attachment_refs_v1'" + ) + + reopened = TeamStore(db_path) + lead = LocalDialect(reopened, None, LEAD, attachments=attachments) + stored = stored_name(ref) + + assert lead.attachment("proj", stored) == (PNG, "image/png") + later = reopened.create_item("proj", LEAD, title="Later", criteria="c") + forged_ref = attachments.put(PNG + b"later", "later.png") + reopened.comment( + "proj", LEAD, later["id"], "ordinary ref", refs=[forged_ref] + ) + reopened.close() + + reopened = TeamStore(db_path) + lead = LocalDialect(reopened, None, LEAD, attachments=attachments) + with pytest.raises(BoardError, match="attachment not found"): + lead.attachment("proj", stored_name(forged_ref)) + reopened.rebuild("proj") + assert lead.attachment("proj", stored) == (PNG, "image/png") + assert reopened.rekey_space("proj", "moved") is True + assert lead.attachment("moved", stored) == (PNG, "image/png") + reopened.close() + + +def test_local_attachment_read_allows_a_directly_linked_item(tmp_path): + from coworker.teams.attachments import stored_name + + lead = local_dialect(tmp_path, actor="lead-1", role="lead") + worker = LocalDialect( + lead.store, lead.journal, NIA, attachments=lead.attachments + ) + mine = lead.create_item("proj", title="Nia task", criteria="c") + lead.assign("proj", mine["id"], "nia") + dependency = lead.create_item("proj", title="Dependency", criteria="c") + lead.assign("proj", dependency["id"], "webb") + ref = lead.attach("proj", dependency["id"], PNG, "dependency.png")["payload"][ + "refs" + ][0] + + with pytest.raises(BoardError, match="attachment not found"): + worker.attachment("proj", stored_name(ref)) + + lead.link("proj", dependency["id"], "blocks", mine["id"]) + assert worker.attachment("proj", stored_name(ref)) == (PNG, "image/png") + + def test_attach_rejects_bad_payloads_over_the_wire(api): client, manager, app = api lead_token = _tokens(manager).mint("lead-1", "lead") @@ -733,6 +1032,75 @@ def test_cli_worker_cannot_show_a_foreign_item(tmp_path, capsys): assert "Private Webb task" not in output.err +def test_cli_worker_cannot_download_a_foreign_attachment(tmp_path, capsys): + from coworker.teams.attachments import stored_name + from coworker.teams.cli import main + + lead = local_dialect(tmp_path, actor="lead-1", role="lead") + foreign = lead.create_item("proj", title="Private Webb task", criteria="c") + lead.assign("proj", foreign["id"], "webb") + ref = lead.attach("proj", foreign["id"], PNG, "private.png")["payload"]["refs"][ + 0 + ] + output = tmp_path / "stolen.png" + + result = main( + [ + "board", + "attachment", + stored_name(ref), + "--out", + str(output), + "--db", + str(tmp_path), + "--space", + "proj", + "--actor", + "nia", + "--role", + "worker", + ] + ) + + assert result == 1 + assert not output.exists() + assert "attachment not found" in capsys.readouterr().err + + +def test_cli_authorized_attachment_download(tmp_path, capsys): + from coworker.teams.attachments import stored_name + from coworker.teams.cli import main + + lead = local_dialect(tmp_path, actor="lead-1", role="lead") + item = lead.create_item("proj", title="Evidence", criteria="c") + ref = lead.attach("proj", item["id"], PNG, "evidence.png")["payload"]["refs"][ + 0 + ] + output = tmp_path / "downloaded.png" + + result = main( + [ + "board", + "attachment", + stored_name(ref), + "--out", + str(output), + "--db", + str(tmp_path), + "--space", + "proj", + "--actor", + "lead-1", + "--role", + "lead", + ] + ) + + assert result == 0 + assert output.read_bytes() == PNG + assert str(output) in capsys.readouterr().out + + def test_cli_token_mint_and_list(tmp_path, capsys): from coworker.teams.cli import main diff --git a/tests/test_team_wake.py b/tests/test_team_wake.py index a3a0bb9d3..5d5e6fd12 100644 --- a/tests/test_team_wake.py +++ b/tests/test_team_wake.py @@ -413,6 +413,91 @@ def test_item_detail_timeline_and_blocker_fact(manager): assert blocked["blocker"] == "need the staging tfvars" +def test_session_attachment_read_is_scoped_to_its_board(manager): + from coworker.sessions import SessionRecord + from coworker.teams import BoardError + from coworker.teams.attachments import stored_name + + space = str(manager.default_workspace) + manager.session_store.save( + SessionRecord( + session_id="sid", + workspace=space, + model="m", + mode="interactive", + messages=[], + agent="cowork", + ) + ) + other = manager.team_store.create_item( + "other-space", LEAD, title="Other board", criteria="c" + ) + ref = manager.attachment_store.put( + b"\x89PNG\r\n\x1a\nprivate", "private.png" + ) + manager.team_store.attach_ref( + "other-space", LEAD, other["id"], "private", ref + ) + + with pytest.raises(BoardError, match="attachment not found"): + manager.board_attachment("sid", stored_name(ref)) + + +def test_session_attachment_route_uses_the_session_board(manager, monkeypatch): + from fastapi.testclient import TestClient + + from coworker.server.app import create_app + from coworker.sessions import SessionRecord + from coworker.teams.attachments import stored_name + + monkeypatch.delenv("COWORKER_API_TOKEN", raising=False) + space = str(manager.default_workspace) + manager.session_store.save( + SessionRecord( + session_id="sid", + workspace=space, + model="m", + mode="interactive", + messages=[], + agent="cowork", + ) + ) + other = manager.team_store.create_item( + "other-space", LEAD, title="Other board", criteria="c" + ) + ref = manager.attachment_store.put( + b"\x89PNG\r\n\x1a\nprivate", "private.png" + ) + manager.team_store.attach_ref( + "other-space", LEAD, other["id"], "private", ref + ) + client = TestClient(create_app(manager)) + + response = client.get( + "/v1/sessions/sid/board/attachment", + params={"name": stored_name(ref)}, + ) + assert response.status_code == 404 + assert response.json() == {"error": "attachment not found"} + + item = manager.team_store.create_item( + space, LEAD, title="This board", criteria="c" + ) + visible = manager.attachment_store.put( + b"\x89PNG\r\n\x1a\nvisible", "visible.png" + ) + manager.team_store.attach_ref(space, LEAD, item["id"], "visible", visible) + allowed = client.get( + "/v1/sessions/sid/board/attachment", + params={"name": stored_name(visible)}, + ) + client.close() + + assert allowed.status_code == 200 + assert allowed.content == b"\x89PNG\r\n\x1a\nvisible" + assert allowed.headers["content-type"] == "image/png" + + def test_feed_interest_follows_the_assignment_relation(store): """Owner ruling 2026-08-17: no per-event addressing — a worker is subscribed to everything on its slice. Send-backs, comment ANSWERS (the silently broken From 9947e7eb282d73a010649fdff1c88f1df09883e5 Mon Sep 17 00:00:00 2001 From: Rakesh Utekar <48244158+rakeshutekar@users.noreply.github.com> Date: Fri, 28 Aug 2026 16:56:40 -0700 Subject: [PATCH 3/3] fix(teams): scope board tokens to one space Bind every external board credential to one opaque space and enforce that scope across HTTP, CLI, MCP, and journal projections. Preserve trusted local cross-board journal behavior while legacy and malformed credentials fail closed. Serialize token registry changes across processes and coordinate board/journal rekeys transactionally so retired credentials cannot recreate an old space. --- coworker/projects.py | 4 +- coworker/server/app.py | 144 ++++++--- coworker/server/manager.py | 16 +- coworker/teams/cli.py | 39 ++- coworker/teams/dialect.py | 4 +- coworker/teams/journal.py | 168 +++++++++- coworker/teams/mcp_server.py | 5 + coworker/teams/store.py | 211 +++++++++---- coworker/teams/tokens.py | 259 +++++++++++++-- tests/test_projects.py | 12 + tests/test_team_journal.py | 99 ++++++ tests/test_team_open_surface.py | 545 ++++++++++++++++++++++++++++++-- 12 files changed, 1313 insertions(+), 193 deletions(-) diff --git a/coworker/projects.py b/coworker/projects.py index b64cd0be1..bc375f475 100644 --- a/coworker/projects.py +++ b/coworker/projects.py @@ -206,7 +206,9 @@ def resolve_board_space( try: team_store.rekey_space(path_key, derived) except Exception: - pass + # A failed cross-store move leaves the path-keyed board authoritative. + # Retry on the next resolution instead of selecting an empty target. + return path_key return derived diff --git a/coworker/server/app.py b/coworker/server/app.py index 457c08110..b8551fdad 100644 --- a/coworker/server/app.py +++ b/coworker/server/app.py @@ -821,28 +821,38 @@ def teams_journal() -> dict[str, Any]: return {"cases": manager.journal_overview()} # ---- The open board surface (OPE-100): token-authenticated `/v1/board` API. - # Identity is the TOKEN (actor+role bound at mint, resolved per request, never - # client-asserted); authority is the STORE — the same double gate in-app agents - # get. This is the one wire protocol every external front door rides: + # Identity and board scope are the TOKEN (actor+role+space bound at mint, + # resolved per request, never client-asserted); object and verb authority are + # the STORE. This is the one wire protocol every external front door rides: # RemoteDialect (the `ocw` CLI, the team-board MCP server, headless instances) # today, a hosted board service later. Tokens are required even on loopback — # they carry identity, not just access. - def _board_actor(request: Request): + def _board_principal(request: Request): auth = request.headers.get("authorization", "") token = auth[7:] if auth.lower().startswith("bearer ") else "" return manager.board_tokens.resolve(token) - def _board(request: Request, handler): - actor = _board_actor(request) - if actor is None: + _INVALID_BOARD_SPACE = object() + + def _board(request: Request, handler, *, requested_space: Any = None): + principal = _board_principal(request) + if principal is None: return JSONResponse( {"error": "board token required (Authorization: Bearer …) — mint" " one with `ocw board token` on the serving machine"}, status_code=401, ) + if requested_space is not None and requested_space != principal.space: + return JSONResponse( + {"error": "board space not found"}, status_code=404 + ) try: - return handler(actor) + # The store lock makes scope validation and the complete handler one + # operation relative to rekey. A request authenticated just before a + # move either finishes before the move or observes the retired scope. + with manager.team_store.authorized_space(principal.space): + return handler(principal.actor, principal.space) except TeamsBoardNotFoundError as error: return JSONResponse({"error": str(error)}, status_code=404) except TeamsAuthorityError as error: @@ -853,12 +863,17 @@ def _board(request: Request, handler): @app.get("/v1/board/whoami") def board_whoami(request: Request): return _board( - request, lambda actor: {"actor": actor.id, "role": actor.role.value} + request, + lambda actor, space: { + "actor": actor.id, + "role": actor.role.value, + "space": space, + }, ) @app.get("/v1/board/spaces") def board_spaces(request: Request): - return _board(request, lambda actor: {"spaces": manager.team_store.spaces()}) + return _board(request, lambda actor, space: {"spaces": [space]}) @app.get("/v1/board/items") def board_list_items( @@ -866,27 +881,33 @@ def board_list_items( ): return _board( request, - lambda actor: { + lambda actor, _: { "items": manager.team_store.list_items( space, actor, state=state or None, assignee=assignee or None ) }, + requested_space=space, ) @app.get("/v1/board/item") def board_get_item(request: Request, space: str, id: int): return _board( request, - lambda actor: manager.team_store.get_item(space, int(id), actor=actor), + lambda actor, _: manager.team_store.get_item( + space, int(id), actor=actor + ), + requested_space=space, ) @app.post("/v1/board/items") def board_create_item(request: Request, body: dict): body = body or {} - def run(actor): + space = str(body.get("space", "")) + + def run(actor, _): item = manager.team_store.create_item( - str(body.get("space", "")), + space, actor, title=str(body.get("title", "")), criteria=str(body.get("criteria", "")), @@ -899,15 +920,17 @@ def run(actor): manager.kick_team_tick() # a new filing is lead-subscription news return item - return _board(request, run) + return _board(request, run, requested_space=space) @app.post("/v1/board/items/transition") def board_transition_item(request: Request, body: dict): body = body or {} - def run(actor): + space = str(body.get("space", "")) + + def run(actor, _): item = manager.team_store.transition( - str(body.get("space", "")), + space, actor, int(body.get("id", 0)), str(body.get("to", "")), @@ -917,29 +940,33 @@ def run(actor): manager.kick_team_tick() # review/blocked should reach the lead now return item - return _board(request, run) + return _board(request, run, requested_space=space) @app.post("/v1/board/items/comment") def board_comment_item(request: Request, body: dict): body = body or {} + space = str(body.get("space", "")) return _board( request, - lambda actor: manager.team_store.comment( - str(body.get("space", "")), + lambda actor, _: manager.team_store.comment( + space, actor, int(body.get("id", 0)), str(body.get("body", "")), refs=[str(ref) for ref in body.get("refs") or []], ), + requested_space=space, ) @app.post("/v1/board/items/assign") def board_assign_item(request: Request, body: dict): body = body or {} - def run(actor): + space = str(body.get("space", "")) + + def run(actor, _): item = manager.team_store.assign( - str(body.get("space", "")), + space, actor, int(body.get("id", 0)), str(body.get("assignee", "")), @@ -947,40 +974,46 @@ def run(actor): manager.kick_team_tick() # the assignee's queue has news return item - return _board(request, run) + return _board(request, run, requested_space=space) @app.post("/v1/board/items/claim") def board_claim_item(request: Request, body: dict): body = body or {} - def run(actor): + space = str(body.get("space", "")) + + def run(actor, _): item = manager.team_store.claim( - str(body.get("space", "")), actor, int(body.get("id", 0)) + space, actor, int(body.get("id", 0)) ) manager.kick_team_tick() # claims land in the lead's feed return item - return _board(request, run) + return _board(request, run, requested_space=space) @app.post("/v1/board/link") def board_link_items(request: Request, body: dict): body = body or {} + space = str(body.get("space", "")) return _board( request, - lambda actor: manager.team_store.link( - str(body.get("space", "")), + lambda actor, _: manager.team_store.link( + space, actor, int(body.get("src", 0)), str(body.get("kind", "")), int(body.get("dst", 0)), ), + requested_space=space, ) @app.post("/v1/board/items/attach") def board_attach(request: Request, body: dict): body = body or {} - def run(actor): + space = str(body.get("space", "")) + + def run(actor, _): raw = str(body.get("data_b64", "")) # Cheap pre-decode bound: base64 is ~4/3 of the payload, so anything # multiples over the cap is refused before allocating the decode. @@ -999,7 +1032,7 @@ def run(actor): ) filename = str(body.get("filename", "")) event = manager.team_store.attach_ref( - str(body.get("space", "")), + space, actor, int(body.get("id", 0)), str(body.get("caption", "")) or f"attached {filename}", @@ -1007,11 +1040,11 @@ def run(actor): ) return {"ref": ref, "seq": event["seq"]} - return _board(request, run) + return _board(request, run, requested_space=space) @app.get("/v1/board/attachment") def board_attachment(request: Request, name: str, space: str): - def run(actor): + def run(actor, _): from fastapi.responses import Response manager.team_store.require_attachment_access(space, actor, name) @@ -1021,20 +1054,26 @@ def run(actor): media_type=manager.attachment_store.mime_for(name), ) - return _board(request, run) + return _board(request, run, requested_space=space) @app.get("/v1/board/policy") def board_get_policy(request: Request, space: str): - return _board(request, lambda actor: manager.team_store.policy(space)) + return _board( + request, + lambda actor, _: manager.team_store.policy(space), + requested_space=space, + ) @app.post("/v1/board/policy") def board_set_policy(request: Request, body: dict): body = body or {} + space = str(body.get("space", "")) return _board( request, - lambda actor: manager.team_store.set_policy( - str(body.get("space", "")), actor, claims=str(body.get("claims", "")) + lambda actor, _: manager.team_store.set_policy( + space, actor, claims=str(body.get("claims", "")) ), + requested_space=space, ) @app.get("/v1/board/pending") @@ -1043,29 +1082,37 @@ def board_pending(request: Request, space: str, limit: int = 200): # follows the assignment relation, same projection in-app workers use. return _board( request, - lambda actor: { + lambda actor, _: { "events": manager.team_store.feed_for( space, actor.id, limit=int(limit) ) }, + requested_space=space, ) @app.post("/v1/board/consume") def board_consume(request: Request, body: dict): body = body or {} - def run(actor): + space = str(body.get("space", "")) + + def run(actor, _): manager.team_store.consume_feed( - str(body.get("space", "")), actor.id, int(body.get("upto_seq", 0)) + space, actor.id, int(body.get("upto_seq", 0)) ) return {"ok": True} - return _board(request, run) + return _board(request, run, requested_space=space) @app.get("/v1/board/journal/cases") def board_journal_cases(request: Request): return _board( - request, lambda actor: {"cases": manager.journal_store.overview(actor)} + request, + lambda actor, space: { + "cases": manager.journal_store.overview( + actor, scope_space=space + ) + }, ) @app.get("/v1/board/journal") @@ -1081,7 +1128,7 @@ def board_journal_read( ): return _board( request, - lambda actor: { + lambda actor, space: { "entries": manager.journal_store.read( actor, case, @@ -1091,6 +1138,7 @@ def board_journal_read( entity=entity or None, include_raw=bool(include_raw), limit=int(limit), + scope_space=space, ) }, ) @@ -1098,18 +1146,26 @@ def board_journal_read( @app.post("/v1/board/journal") def board_journal_append(request: Request, body: dict): body = body or {} + supplied_space = body.get("space") + if "space" not in body: + requested_space = None + elif isinstance(supplied_space, str): + requested_space = supplied_space + else: + requested_space = _INVALID_BOARD_SPACE return _board( request, - lambda actor: manager.journal_store.append( + lambda actor, space: manager.journal_store.append( actor, str(body.get("case", "")), str(body.get("body", "")), kind=str(body.get("kind") or "note"), - space=str(body.get("space") or "") or None, + space=space, item=int(body["item"]) if body.get("item") is not None else None, entities=[str(e) for e in body.get("entities") or []], refs=[str(ref) for ref in body.get("refs") or []], ), + requested_space=requested_space, ) @app.get("/v1/memory") diff --git a/coworker/server/manager.py b/coworker/server/manager.py index 130ede1f2..a09a508a7 100644 --- a/coworker/server/manager.py +++ b/coworker/server/manager.py @@ -282,13 +282,21 @@ def __init__( # and assignment feeds journal-case grants. Verbs register per-session behind # the persona's `team:` trait; the registry holds rosters (lead/worker # sessions per board) that the wake plumbing walks. + self.board_tokens = BoardTokens(base / "board-tokens.json") self.journal_store = JournalStore(base / "journal.db") - self.team_store = TeamStore(base / "teams.db", journal=self.journal_store) + self.team_store = TeamStore( + base / "teams.db", + journal=self.journal_store, + space_rekeys=self.board_tokens, + ) + # A persisted token tombstone is the rekey intent. If a process exited + # after recording it but before token cleanup, replay the idempotent move. + for old_space, new_space in self.board_tokens.space_rekeys().items(): + self.team_store.rekey_space(old_space, new_space) self.chat_store = ChatStore(base / "chat.db") self.teams = TeamRegistry(base / "teams.json") - # External board clients (OPE-100): join tokens bind actor+role; the - # `/v1/board` API resolves them and the store enforces authority. - self.board_tokens = BoardTokens(base / "board-tokens.json") + # External board clients (OPE-100): join tokens bind actor+role+space; + # `/v1/board` enforces space scope and the store enforces object authority. # Work-item attachments (OPE-105): content-addressed blobs next to the # board; the log carries only `attachment://` refs. self.attachment_store = AttachmentStore(base / "attachments") diff --git a/coworker/teams/cli.py b/coworker/teams/cli.py index 98d730825..a8382ca3f 100644 --- a/coworker/teams/cli.py +++ b/coworker/teams/cli.py @@ -20,6 +20,7 @@ from __future__ import annotations import argparse +import hashlib import json import os import sys @@ -126,6 +127,7 @@ def cmd(name: str, func, help: str, parent=board_sub): p.add_argument( "--role", choices=("worker", "lead", "user"), default="worker" ) + p.add_argument("--space", default="", help="exact board space the token may use") p.add_argument("--label", default="", help="what this token is for (mint)") p.add_argument("--prefix", default="", help="token prefix to revoke") p.add_argument("--db", default="", help="state dir holding the registry") @@ -195,7 +197,7 @@ def _dialect(args): return local_dialect(args.db, actor=args.local_actor, role=args.local_role) server = _discover_server() if server is not None: - return RemoteDialect(server, _local_cli_token()) + return RemoteDialect(server, _local_cli_token(_space(args))) from ..secrets import state_dir return local_dialect(state_dir(), actor=args.local_actor, role=args.local_role) @@ -226,7 +228,7 @@ def _discover_server() -> Optional[str]: return None -def _local_cli_token() -> str: +def _local_cli_token(space: str) -> str: """The CLI's own user token against the local server. Minted once into the shared registry; the plaintext is cached user-only in the state dir — the user's own credential on the user's own machine, same pattern as the sidecar @@ -235,15 +237,17 @@ def _local_cli_token() -> str: from .tokens import BoardTokens - cache = state_dir() / "ocw-cli.token" + scope_key = hashlib.sha256(space.encode("utf-8")).hexdigest()[:16] + cache = state_dir() / f"ocw-cli-{scope_key}.token" tokens = BoardTokens(state_dir() / "board-tokens.json") try: cached = cache.read_text().strip() - if cached and tokens.resolve(cached) is not None: + principal = tokens.resolve(cached) if cached else None + if principal is not None and principal.space == space: return cached except OSError: pass - token = tokens.mint("user", "user", label="local ocw CLI") + token = tokens.mint("user", "user", space=space, label="local ocw CLI") write_private_text(cache, token + "\n") return token @@ -416,10 +420,25 @@ def _cmd_token(args) -> int: if not args.actor: print("error: --actor is required to mint", file=sys.stderr) return 1 - token = tokens.mint(args.actor, args.role, label=args.label) + if not args.space: + print("error: --space is required to mint", file=sys.stderr) + return 1 + space = args.space + checked_space = space.strip() + if not checked_space or checked_space == "*": + print("error: one exact board space is required", file=sys.stderr) + return 1 + try: + token = tokens.mint( + args.actor, args.role, space=space, label=args.label + ) + except ValueError as error: + print(f"error: {error}", file=sys.stderr) + return 1 print(token) print( - f"# binds actor '{args.actor}' as {args.role}; shown once — store it" + f"# binds actor '{args.actor}' as {args.role} in space {space!r};" + " shown once — store it" " in the client's config (OCW_BOARD_TOKEN)", file=sys.stderr, ) @@ -434,7 +453,11 @@ def _cmd_token(args) -> int: return 0 for entry in entries: label = f" ({entry['label']})" if entry["label"] else "" - print(f"{entry['prefix']}… {entry['actor']:<16} {entry['role']:<8}{label}") + space = entry.get("space") or "unscoped (disabled)" + print( + f"{entry['prefix']}… {entry['actor']:<16} {entry['role']:<8}" + f" {space}{label}" + ) if not entries: print("no tokens") return 0 diff --git a/coworker/teams/dialect.py b/coworker/teams/dialect.py index df66ac6e6..d7b80ec32 100644 --- a/coworker/teams/dialect.py +++ b/coworker/teams/dialect.py @@ -303,8 +303,8 @@ def _need_journal(self) -> None: class RemoteDialect: """The `/v1/board` HTTP client. `base_url` is an OpenWorker sidecar or a hosted - board service; the Bearer token carries identity — the server resolves it to an - actor+role, so this client never states who it is, it proves it.""" + board service; the Bearer token carries identity and one board-space scope, so + this client never states its own authority — the server resolves it.""" def __init__( self, base_url: str, token: str, *, client: Any = None, timeout: float = 30.0 diff --git a/coworker/teams/journal.py b/coworker/teams/journal.py index 1cdc180d8..d436abae7 100644 --- a/coworker/teams/journal.py +++ b/coworker/teams/journal.py @@ -183,7 +183,13 @@ def append( (case, record["hash"], ts), ) # A new case belongs to whoever opened it. - self._grant_locked(case, actor.id, source="creator") + self._grant_locked( + case, + actor.id, + source="creator", + space=space or "", + item_id=item, + ) else: self._conn.execute( "UPDATE journal_meta SET head_hash = ? WHERE case_id = ?", @@ -207,6 +213,7 @@ def read( since_seq: int = 0, include_raw: bool = False, limit: int = 100, + scope_space: Optional[str] = None, ) -> list[dict[str, Any]]: """Filtered read. `raw` captures are skipped unless asked for (by `kind="raw"` or `include_raw`) so dumps never bury the signal entries.""" @@ -217,6 +224,9 @@ def read( if item is not None: where.append("item_id = ?") params.append(item) + if scope_space is not None: + where.append("space = ?") + params.append(scope_space) if author: where.append("actor = ?") params.append(author) @@ -244,16 +254,25 @@ def read( break return out - def overview(self, actor: Actor) -> list[dict[str, Any]]: + def overview( + self, actor: Actor, *, scope_space: Optional[str] = None + ) -> list[dict[str, Any]]: """Case list with entry counts and last activity — the rail's summary view.""" - visible = self.cases(actor) + visible = self.cases(actor, scope_space=scope_space) if not visible: return [] with self._lock: - rows = self._conn.execute( - "SELECT case_id, COUNT(*) AS entries, MAX(ts) AS last_ts" - " FROM journal_entries GROUP BY case_id" - ).fetchall() + if scope_space is None: + rows = self._conn.execute( + "SELECT case_id, COUNT(*) AS entries, MAX(ts) AS last_ts" + " FROM journal_entries GROUP BY case_id" + ).fetchall() + else: + rows = self._conn.execute( + "SELECT case_id, COUNT(*) AS entries, MAX(ts) AS last_ts" + " FROM journal_entries WHERE space = ? GROUP BY case_id", + (scope_space,), + ).fetchall() counts = {row["case_id"]: dict(row) for row in rows} return [ { @@ -264,7 +283,9 @@ def overview(self, actor: Actor) -> list[dict[str, Any]]: for case in visible ] - def cases(self, actor: Actor) -> list[str]: + def cases( + self, actor: Actor, *, scope_space: Optional[str] = None + ) -> list[str]: """Cases visible to this actor (all of them for the user).""" with self._lock: if actor.role == Role.USER: @@ -277,7 +298,16 @@ def cases(self, actor: Actor) -> list[str]: " ORDER BY case_id", (actor.id,), ).fetchall() - return [row["case_id"] for row in rows] + visible = [row["case_id"] for row in rows] + if scope_space is None or not visible: + return visible + scoped_rows = self._conn.execute( + "SELECT case_id FROM journal_entries WHERE space = ?" + " UNION SELECT case_id FROM journal_grants WHERE space = ?", + (scope_space, scope_space), + ).fetchall() + scoped = {row["case_id"] for row in scoped_rows} + return [case for case in visible if case in scoped] # ---------------------------------------------------------------------- grants @@ -308,7 +338,14 @@ def revoke(self, actor: Actor, case: str, principal: str) -> None: ) self._conn.commit() - def ensure_case(self, case: str, creator: str) -> None: + def ensure_case( + self, + case: str, + creator: str, + *, + space: str = "", + item_id: Optional[int] = None, + ) -> None: """Create a case (empty, chain at genesis) if it doesn't exist, granting its creator. Called by the board when an item attaches a case ref — so the case belongs to whoever attached it, not to whichever assignee @@ -323,7 +360,29 @@ def ensure_case(self, case: str, creator: str) -> None: " VALUES (?, ?, ?)", (case, GENESIS, datetime.now(timezone.utc).isoformat()), ) - self._grant_locked(case, creator, source="creator") + self._grant_locked( + case, + creator, + source="creator", + space=space, + item_id=item_id, + ) + self._conn.commit() + return + grants = self._conn.execute( + "SELECT source FROM journal_grants" + " WHERE case_id = ? AND principal = ?", + (case, creator), + ).fetchall() + for grant in grants: + self._grant_locked( + case, + creator, + source=grant["source"], + space=space, + item_id=item_id, + ) + if grants: self._conn.commit() def sync_assignment( @@ -373,6 +432,93 @@ def verify_chain(self, case: str) -> int: raise ChainError("case log ends before the recorded head — tail deleted") return len(rows) + def rekey_space(self, old: str, new: str) -> None: + """Move journal projections with a board-space identity migration. + + Cases may span spaces, so changing even one entry requires recomputing + that case's complete hash chain in sequence order. + """ + if old == new: + return + with self._lock: + try: + self._rekey_space_locked(self._conn, old, new) + self._conn.commit() + except Exception: + self._conn.rollback() + raise + + def _rekey_space_locked( + self, + connection: sqlite3.Connection, + old: str, + new: str, + *, + schema: str = "main", + ) -> None: + """Apply a rekey using an already-locked caller-owned transaction.""" + if schema not in {"main", "journal_rekey"}: + raise ValueError(f"unsupported journal schema: {schema}") + entries = f"{schema}.journal_entries" + grants_table = f"{schema}.journal_grants" + meta = f"{schema}.journal_meta" + cases = [ + row["case_id"] + for row in connection.execute( + f"SELECT DISTINCT case_id FROM {entries} WHERE space = ?", + (old,), + ).fetchall() + ] + grants = connection.execute( + f"SELECT case_id, principal, source, item_id FROM {grants_table}" + " WHERE space = ?", + (old,), + ).fetchall() + connection.execute( + f"UPDATE {entries} SET space = ? WHERE space = ?", (new, old) + ) + for grant in grants: + connection.execute( + f"INSERT OR IGNORE INTO {grants_table}" + " (case_id, principal, source, space, item_id)" + " VALUES (?, ?, ?, ?, ?)", + ( + grant["case_id"], + grant["principal"], + grant["source"], + new, + grant["item_id"], + ), + ) + connection.execute( + f"DELETE FROM {grants_table} WHERE space = ?", (old,) + ) + for case in cases: + rows = connection.execute( + f"SELECT * FROM {entries} WHERE case_id = ? ORDER BY seq", + (case,), + ).fetchall() + prev = GENESIS + for row in rows: + record = { + key: (prev if key == "prev_hash" else row[key]) + for key in _HASHED_FIELDS + } + digest = _hash(record, fields=_HASHED_FIELDS) + connection.execute( + f"UPDATE {entries} SET prev_hash = ?, hash = ? WHERE seq = ?", + (prev, digest, row["seq"]), + ) + prev = digest + connection.execute( + f"UPDATE {meta} SET head_hash = ? WHERE case_id = ?", + (prev, case), + ) + + def rekey_lock(self): + """Return the lock used by the board's coordinated two-database rekey.""" + return self._lock + def close(self) -> None: self._conn.close() diff --git a/coworker/teams/mcp_server.py b/coworker/teams/mcp_server.py index 1f9466a8d..aa3d692f5 100644 --- a/coworker/teams/mcp_server.py +++ b/coworker/teams/mcp_server.py @@ -23,6 +23,11 @@ def build(dialect, *, space: str): from mcp.server.fastmcp import FastMCP who = dialect.whoami() + token_space = who.get("space") + if token_space is not None and token_space != space: + raise BoardError( + f"board token is scoped to {token_space!r}, not MCP space {space!r}" + ) role = who.get("role", "worker") mcp = FastMCP( "team-board", diff --git a/coworker/teams/store.py b/coworker/teams/store.py index 188a9ee01..ec779efb4 100644 --- a/coworker/teams/store.py +++ b/coworker/teams/store.py @@ -27,6 +27,7 @@ import json import sqlite3 import threading +from contextlib import contextmanager, nullcontext from datetime import datetime, timezone from pathlib import Path from typing import Any, Optional @@ -79,11 +80,18 @@ class TeamStore: - def __init__(self, db_path: str | Path, *, journal: Any = None) -> None: + def __init__( + self, + db_path: str | Path, + *, + journal: Any = None, + space_rekeys: Any = None, + ) -> None: # `journal` is a teams.journal.JournalStore when wired: assignment feeds # case grants ("sharing rides assignment"). Optional so the board works # standalone (tests, boards with no journal). self.journal = journal + self._space_rekeys = space_rekeys self.db_path = str(db_path) if self.db_path != ":memory:": Path(self.db_path).expanduser().parent.mkdir(parents=True, exist_ok=True) @@ -529,7 +537,9 @@ def create_item( }, ) if self.journal is not None and case: - self.journal.ensure_case(case, actor.id) + self.journal.ensure_case( + case, actor.id, space=space, item_id=item_id + ) return self.get_item(space, item_id, actor=actor, seq=event["seq"]) def list_items( @@ -1165,6 +1175,16 @@ def event_count(self, space: str) -> int: ).fetchone() return int(row[0]) if row else 0 + @contextmanager + def authorized_space(self, space: str): + """Serialize a remote operation with rekey and reject retired scopes.""" + with self._lock: + if self._space_rekeys is not None and not self._space_rekeys.is_space_active( + space + ): + raise BoardNotFoundError("board space not found") + yield + def rekey_space(self, old: str, new: str) -> bool: """Move one space's records under a new key — the twentieth-pass one-time path→git migration. `space` participates in the hash chain, so the chain @@ -1174,80 +1194,139 @@ def rekey_space(self, old: str, new: str) -> bool: if old == new: return True with self._lock: + has_old = self._conn.execute( + "SELECT 1 FROM team_events WHERE space = ? LIMIT 1", (old,) + ).fetchone() has_new = self._conn.execute( "SELECT 1 FROM team_events WHERE space = ? LIMIT 1", (new,) ).fetchone() if has_new: + pending = ( + self._space_rekeys.space_rekeys().get(old) == new + if self._space_rekeys is not None + else False + ) + if pending and not has_old: + # The databases committed before token cleanup (for example, + # a process exit at that boundary). The move is complete. + self._finish_space_rekey(old, new) + return True return False rows = self._conn.execute( "SELECT * FROM team_events WHERE space = ? ORDER BY seq", (old,) ).fetchall() - try: - prev = GENESIS - for row in rows: - record = { - "ts": row["ts"], - "space": new, - "kind": row["kind"], - "actor": row["actor"], - "actor_role": row["actor_role"], - "item_id": row["item_id"], - "case_id": row["case_id"], - "recipient": row["recipient"], - "payload": row["payload"], - "taint": row["taint"], - "prev_hash": prev, - } - record["hash"] = _hash(record) - self._conn.execute( - "UPDATE team_events SET space = ?, prev_hash = ?, hash = ? " - "WHERE seq = ?", - (new, prev, record["hash"], row["seq"]), - ) - prev = record["hash"] - for table in ( - "team_items", - "team_links", - "team_attachment_refs", - "team_settings", - ): - self._conn.execute( - f"UPDATE {table} SET space = ? WHERE space = ?", (new, old) - ) - # Cursor keys embed the space as a suffix ("feed::", - # "sub::") — rewrite the suffix, keep consumed positions. - cur_rows = self._conn.execute( - "SELECT cursor_key FROM team_cursors WHERE cursor_key LIKE ?", - ("%:" + old,), - ).fetchall() - for crow in cur_rows: - new_key = crow["cursor_key"][: -len(old)] + new - self._conn.execute( - "UPDATE OR REPLACE team_cursors SET cursor_key = ? " - "WHERE cursor_key = ?", - (new_key, crow["cursor_key"]), - ) - meta = self._conn.execute( - "SELECT watermark FROM team_meta WHERE space = ?", (old,) - ).fetchone() - if meta is not None and rows: - self._conn.execute( - "DELETE FROM team_meta WHERE space = ?", (old,) - ) - self._conn.execute( - "INSERT INTO team_meta (space, head_hash, watermark) " - "VALUES (?, ?, ?) ON CONFLICT(space) DO UPDATE SET " - "head_hash = excluded.head_hash, watermark = excluded.watermark", - (new, prev, meta["watermark"]), - ) - elif meta is not None: - self._conn.execute("DELETE FROM team_meta WHERE space = ?", (old,)) - except Exception: - self._conn.rollback() - raise - self._conn.commit() + journal_lock = ( + self.journal.rekey_lock() if self.journal is not None else nullcontext() + ) + with journal_lock: + attached = False + created_rekey_intent = False + try: + if self.journal is not None: + if self.journal.db_path == ":memory:": + raise BoardError( + "coordinated space rekey requires a file-backed journal" + ) + self._conn.execute( + "ATTACH DATABASE ? AS journal_rekey", + (self.journal.db_path,), + ) + attached = True + if self._space_rekeys is not None: + created_rekey_intent = ( + self._space_rekeys.begin_space_rekey(old, new) + ) + self._conn.execute("BEGIN IMMEDIATE") + prev = GENESIS + for row in rows: + record = { + "ts": row["ts"], + "space": new, + "kind": row["kind"], + "actor": row["actor"], + "actor_role": row["actor_role"], + "item_id": row["item_id"], + "case_id": row["case_id"], + "recipient": row["recipient"], + "payload": row["payload"], + "taint": row["taint"], + "prev_hash": prev, + } + record["hash"] = _hash(record) + self._conn.execute( + "UPDATE team_events SET space = ?, prev_hash = ?, hash = ? " + "WHERE seq = ?", + (new, prev, record["hash"], row["seq"]), + ) + prev = record["hash"] + for table in ( + "team_items", + "team_links", + "team_attachment_refs", + "team_settings", + ): + self._conn.execute( + f"UPDATE {table} SET space = ? WHERE space = ?", (new, old) + ) + # Cursor keys embed the space as a suffix + # ("feed::", "sub::"). + cur_rows = self._conn.execute( + "SELECT cursor_key FROM team_cursors WHERE cursor_key LIKE ?", + ("%:" + old,), + ).fetchall() + for crow in cur_rows: + new_key = crow["cursor_key"][: -len(old)] + new + self._conn.execute( + "UPDATE OR REPLACE team_cursors SET cursor_key = ? " + "WHERE cursor_key = ?", + (new_key, crow["cursor_key"]), + ) + meta = self._conn.execute( + "SELECT watermark FROM team_meta WHERE space = ?", (old,) + ).fetchone() + if meta is not None and rows: + self._conn.execute( + "DELETE FROM team_meta WHERE space = ?", (old,) + ) + self._conn.execute( + "INSERT INTO team_meta (space, head_hash, watermark) " + "VALUES (?, ?, ?) ON CONFLICT(space) DO UPDATE SET " + "head_hash = excluded.head_hash, watermark = excluded.watermark", + (new, prev, meta["watermark"]), + ) + elif meta is not None: + self._conn.execute( + "DELETE FROM team_meta WHERE space = ?", (old,) + ) + if self.journal is not None: + self.journal._rekey_space_locked( + self._conn, old, new, schema="journal_rekey" + ) + self._commit_rekey() + except Exception: + self._conn.rollback() + if created_rekey_intent: + self._space_rekeys.cancel_space_rekey(old, new) + raise + finally: + if attached: + self._conn.execute("DETACH DATABASE journal_rekey") + self._finish_space_rekey(old, new) return True + def _commit_rekey(self) -> None: + self._conn.commit() + + def _finish_space_rekey(self, old: str, new: str) -> None: + if self._space_rekeys is None: + return + try: + self._space_rekeys.finish_space_rekey(old, new) + except OSError: + # The persisted tombstone already denies old credentials. Startup + # recovery retries this token-file cleanup without risking the move. + pass + def _head_hash(self, space: str) -> str: row = self._conn.execute( "SELECT head_hash FROM team_meta WHERE space = ?", (space,) diff --git a/coworker/teams/tokens.py b/coworker/teams/tokens.py index 32740848b..8373bb551 100644 --- a/coworker/teams/tokens.py +++ b/coworker/teams/tokens.py @@ -1,10 +1,10 @@ """Board join tokens — identity for external board clients. -A token binds an ACTOR and a ROLE server-side: an external harness (another agent -CLI, a headless OpenWorker, the `ocw` CLI from a second machine) presents the token -and the server resolves who it is — the client never states its own identity, and a -worker token cannot claim to be the lead. Authority then falls to the store, same -as for in-app agents: the token is identity, the store is the gate. +A token binds an ACTOR, ROLE, and exactly one board SPACE server-side: an external +harness (another agent CLI, a headless OpenWorker, the `ocw` CLI from a second +machine) presents the token and the server resolves both who it is and which board +it may address. The client never states its own identity or authority. Space scope +is the transport gate; actor visibility and verb authority remain store gates. Storage is hash-only (sha256): the plaintext is shown once at mint and never persisted, so the registry file leaking doesn't leak the credentials. Revocation is @@ -15,15 +15,29 @@ import hashlib import json +import os import secrets import threading +from contextlib import contextmanager +from dataclasses import dataclass from datetime import datetime, timezone from pathlib import Path from typing import Any, Optional +from ..secrets import write_private_text from .model import Actor, Role _TOKEN_PREFIX = "owb_" # OpenWorker board — greppable in configs, meaningless to guess +_SPACE_REKEYS_KEY = "__space_rekeys__" +_INVALID_REGISTRY_KEY = "__invalid_registry__" + + +@dataclass(frozen=True) +class BoardPrincipal: + """The identity and one board-space capability carried by a token.""" + + actor: Actor + space: str class BoardTokens: @@ -31,18 +45,31 @@ def __init__(self, path: str | Path) -> None: self.path = Path(path).expanduser() self._lock = threading.Lock() - def mint(self, actor: str, role: str = "worker", *, label: str = "") -> str: - """Create a token for one actor identity; returns the plaintext ONCE.""" + def mint( + self, actor: str, role: str = "worker", *, space: str, label: str = "" + ) -> str: + """Create a token for one actor in one board space; return plaintext once.""" actor = (actor or "").strip() if not actor: raise ValueError("actor is required") - Role(role) # validate early — a bad role should fail at mint, not at use + resolved_role = Role(role) + if resolved_role == Role.SYSTEM: + raise ValueError("system actors cannot use board tokens") + if not isinstance(space, str): + raise ValueError("one exact board space is required") + checked_space = space.strip() + if not checked_space or checked_space == "*": + raise ValueError("one exact board space is required") token = _TOKEN_PREFIX + secrets.token_urlsafe(32) - with self._lock: + with self._locked(): entries = self._load() + self._require_valid(entries) + if space in self._space_rekeys(entries): + raise ValueError(f"board space {space!r} is retired after rekey") entries[_digest(token)] = { "actor": actor, - "role": role, + "role": resolved_role.value, + "space": space, "label": label, "prefix": token[:12], "created_ts": datetime.now(timezone.utc).isoformat(), @@ -50,48 +77,228 @@ def mint(self, actor: str, role: str = "worker", *, label: str = "") -> str: self._save(entries) return token - def resolve(self, token: str) -> Optional[Actor]: + def resolve(self, token: str) -> Optional[BoardPrincipal]: if not token: return None - with self._lock: - entry = self._load().get(_digest(token)) - if entry is None: + with self._locked(): + entries = self._load() + valid = self._registry_valid(entries) + entry = entries.get(_digest(token)) + retired = self._space_rekeys(entries) + if not valid: + return None + if not isinstance(entry, dict): + return None + raw_space = entry.get("space") + if not isinstance(raw_space, str): + return None + checked_space = raw_space.strip() + if not checked_space or checked_space == "*": + return None + if raw_space in retired: + return None + try: + raw_actor = entry["actor"] + if not isinstance(raw_actor, str): + return None + actor_id = raw_actor.strip() + if not actor_id: + return None + role = Role(entry["role"]) + if role == Role.SYSTEM: + return None + actor = Actor(id=actor_id, role=role) + except (KeyError, TypeError, ValueError): return None - return Actor(id=entry["actor"], role=Role(entry["role"])) + return BoardPrincipal(actor=actor, space=raw_space) def entries(self) -> list[dict[str, Any]]: - with self._lock: - return sorted(self._load().values(), key=lambda e: e["created_ts"]) + with self._locked(): + loaded = self._load() + if not self._registry_valid(loaded): + return [] + entries = [ + entry + for key, entry in loaded.items() + if key != _SPACE_REKEYS_KEY and isinstance(entry, dict) + ] + return sorted(entries, key=lambda entry: str(entry.get("created_ts", ""))) def revoke(self, prefix: str) -> int: """Revoke every token whose display prefix matches; returns the count.""" prefix = (prefix or "").strip() if not prefix: return 0 - with self._lock: + with self._locked(): entries = self._load() + if not self._registry_valid(entries): + return 0 keep = { key: entry for key, entry in entries.items() - if not entry["prefix"].startswith(prefix) + if key == _SPACE_REKEYS_KEY + or not isinstance(entry, dict) + or not str(entry.get("prefix", "")).startswith(prefix) } removed = len(entries) - len(keep) if removed: self._save(keep) return removed - def _load(self) -> dict[str, dict[str, Any]]: + def begin_space_rekey(self, old: str, new: str) -> bool: + """Persist a fail-closed rekey intent while leaving tokens recoverable.""" + with self._locked(): + entries = self._load() + self._require_valid(entries) + rekeys = self._space_rekeys(entries) + existing = rekeys.get(old) + if existing is not None and existing != new: + raise ValueError( + f"board space {old!r} is already retired to {existing!r}" + ) + if existing == new: + return False + rekeys[old] = new + entries[_SPACE_REKEYS_KEY] = rekeys + self._save(entries) + return True + + def cancel_space_rekey(self, old: str, new: str) -> None: + """Reactivate tokens when the coordinated database transaction rolls back.""" + with self._locked(): + entries = self._load() + self._require_valid(entries) + rekeys = self._space_rekeys(entries) + if rekeys.get(old) != new: + return + del rekeys[old] + if rekeys: + entries[_SPACE_REKEYS_KEY] = rekeys + else: + entries.pop(_SPACE_REKEYS_KEY, None) + self._save(entries) + + def finish_space_rekey(self, old: str, new: str) -> int: + """Revoke old credentials but retain the tombstone against future minting.""" + with self._locked(): + entries = self._load() + self._require_valid(entries) + rekeys = self._space_rekeys(entries) + if rekeys.get(old) != new: + raise ValueError(f"no pending board-space rekey {old!r} to {new!r}") + keep = { + key: entry + for key, entry in entries.items() + if key == _SPACE_REKEYS_KEY + or not isinstance(entry, dict) + or entry.get("space") != old + } + removed = len(entries) - len(keep) + self._save(keep) + return removed + + def space_rekeys(self) -> dict[str, str]: + """Return persisted rekey intents/tombstones for startup recovery.""" + with self._locked(): + entries = self._load() + if not self._registry_valid(entries): + return {} + return dict(self._space_rekeys(entries)) + + def is_space_active(self, space: str) -> bool: + with self._locked(): + entries = self._load() + return self._registry_valid(entries) and space not in self._space_rekeys( + entries + ) + + @contextmanager + def _locked(self): + """Serialize registry read/modify/write across threads and processes.""" + with self._lock: + with _interprocess_lock(self.path.with_name(self.path.name + ".lock")): + yield + + def _load(self) -> dict[str, Any]: try: - return json.loads(self.path.read_text()) + entries = json.loads(self.path.read_text()) + except FileNotFoundError: + return {} except (OSError, ValueError): + return {_INVALID_REGISTRY_KEY: True} + return entries if isinstance(entries, dict) else {_INVALID_REGISTRY_KEY: True} + + @staticmethod + def _registry_valid(entries: dict[str, Any]) -> bool: + if _INVALID_REGISTRY_KEY in entries: + return False + raw = entries.get(_SPACE_REKEYS_KEY, {}) + return isinstance(raw, dict) and all( + isinstance(old, str) and isinstance(new, str) + for old, new in raw.items() + ) + + @classmethod + def _require_valid(cls, entries: dict[str, Any]) -> None: + if not cls._registry_valid(entries): + raise ValueError("board token registry is malformed; repair it before use") + + @staticmethod + def _space_rekeys(entries: dict[str, Any]) -> dict[str, str]: + raw = entries.get(_SPACE_REKEYS_KEY) + if not isinstance(raw, dict): return {} + return { + old: new + for old, new in raw.items() + if isinstance(old, str) and isinstance(new, str) + } - def _save(self, entries: dict[str, dict[str, Any]]) -> None: - self.path.parent.mkdir(parents=True, exist_ok=True) - tmp = self.path.with_suffix(".tmp") - tmp.write_text(json.dumps(entries, indent=2)) - tmp.replace(self.path) + def _save(self, entries: dict[str, Any]) -> None: + write_private_text(self.path, json.dumps(entries, indent=2)) def _digest(token: str) -> str: return hashlib.sha256(token.encode("utf-8")).hexdigest() + + +@contextmanager +def _interprocess_lock(path: Path): + """Exclusive advisory lock shared by the server and administrative CLIs.""" + path.parent.mkdir(parents=True, exist_ok=True) + flags = os.O_RDWR | os.O_CREAT + if hasattr(os, "O_CLOEXEC"): + flags |= os.O_CLOEXEC + if hasattr(os, "O_NOFOLLOW"): + flags |= os.O_NOFOLLOW + fd = os.open(path, flags, 0o600) + try: + try: + os.chmod(path, 0o600) + except OSError: + pass + if os.name == "nt": + import msvcrt + + if os.fstat(fd).st_size == 0: + os.write(fd, b"\0") + os.lseek(fd, 0, os.SEEK_SET) + msvcrt.locking(fd, msvcrt.LK_LOCK, 1) + else: + import fcntl + + fcntl.flock(fd, fcntl.LOCK_EX) + yield + finally: + try: + if os.name == "nt": + import msvcrt + + os.lseek(fd, 0, os.SEEK_SET) + msvcrt.locking(fd, msvcrt.LK_UNLCK, 1) + else: + import fcntl + + fcntl.flock(fd, fcntl.LOCK_UN) + finally: + os.close(fd) diff --git a/tests/test_projects.py b/tests/test_projects.py index dc2b0e726..33a56c1c8 100644 --- a/tests/test_projects.py +++ b/tests/test_projects.py @@ -162,6 +162,18 @@ def test_derivation_migrates_memory(self, repo, tmp_path): assert key == str(repo.resolve()) assert [m.content for m in store.list(workspace=key)] == ["learned here"] + def test_board_derivation_keeps_old_key_when_rekey_fails(self, repo): + sub = repo / "src" + sub.mkdir() + + class FailingStore: + def rekey_space(self, old, new): + raise RuntimeError("injected rekey failure") + + assert resolve_board_space( + str(sub), team_store=FailingStore() + ) == str(sub.resolve()) + def test_presence_counts(self, tmp_path): mstore = SQLiteMemoryStore(tmp_path / "m.db") tstore = TeamStore(tmp_path / "b.db") diff --git a/tests/test_team_journal.py b/tests/test_team_journal.py index 7d18efe93..befcf78c4 100644 --- a/tests/test_team_journal.py +++ b/tests/test_team_journal.py @@ -13,6 +13,7 @@ TeamStore, ) from coworker.teams.model import JOURNAL_BODY_LIMIT +from coworker.teams.tokens import BoardTokens USER = Actor(id="user", role=Role.USER) LEAD = Actor(id="lead-1", role=Role.LEAD) @@ -80,6 +81,104 @@ def test_cases_span_boards_and_teams(board, journal): assert "loose-threads" in journal.cases(USER) +def test_scoped_journal_reads_project_cross_board_cases_to_one_space(journal): + alpha = journal.append(LEAD, "ops", "alpha finding", space="alpha", item=1) + journal.append(LEAD, "ops", "beta finding", space="beta", item=1) + journal.append(LEAD, "beta-only", "other board", space="beta", item=2) + + assert [entry["body"] for entry in journal.read( + LEAD, "ops", scope_space="alpha" + )] == ["alpha finding"] + assert journal.cases(LEAD, scope_space="alpha") == ["ops"] + assert journal.overview(LEAD, scope_space="alpha") == [ + {"case": "ops", "entries": 1, "last_ts": alpha["ts"]} + ] + + # The trusted in-process view remains case-global. + assert len(journal.read(LEAD, "ops")) == 2 + assert journal.cases(LEAD) == ["beta-only", "ops"] + + +def test_board_rekey_moves_journal_entries_and_grants(board, journal): + item_id = case_item(board, case="ops", assignee="worker-1", space="old") + journal.append( + WORKER, "ops", "old-space entry", space="old", item=item_id + ) + journal.append(LEAD, "ops", "other-space entry", space="other", item=7) + + assert board.rekey_space("old", "new") is True + + assert journal.read(WORKER, "ops", scope_space="old") == [] + assert [ + entry["body"] + for entry in journal.read(WORKER, "ops", scope_space="new") + ] == ["old-space entry"] + assert journal.cases(WORKER, scope_space="new") == ["ops"] + assert journal.verify_chain("ops") == 2 + + +def test_board_and_journal_rekey_roll_back_together(tmp_path, monkeypatch): + journal = JournalStore(tmp_path / "journal.db") + tokens = BoardTokens(tmp_path / "board-tokens.json") + board = TeamStore( + tmp_path / "teams.db", journal=journal, space_rekeys=tokens + ) + token = tokens.mint("lead-1", "lead", space="old") + item_id = case_item(board, case="atomic", space="old") + journal.append(LEAD, "atomic", "before", space="old", item=item_id) + + def fail_commit(): + raise RuntimeError("injected commit failure") + + monkeypatch.setattr(board, "_commit_rekey", fail_commit) + with pytest.raises(RuntimeError, match="injected"): + board.rekey_space("old", "new") + + assert board.event_count("old") > 0 + assert board.event_count("new") == 0 + assert [entry["body"] for entry in journal.read(LEAD, "atomic")] == ["before"] + assert journal.read(LEAD, "atomic", scope_space="new") == [] + assert tokens.resolve(token) is not None + + board.close() + journal.close() + + +def test_rekey_recovery_failure_keeps_an_existing_tombstone(tmp_path, monkeypatch): + journal = JournalStore(tmp_path / "journal.db") + tokens = BoardTokens(tmp_path / "board-tokens.json") + board = TeamStore( + tmp_path / "teams.db", journal=journal, space_rekeys=tokens + ) + board.create_item("old", LEAD, title="Pending", criteria="moves") + assert tokens.begin_space_rekey("old", "new") is True + + def fail_commit(): + raise RuntimeError("injected recovery failure") + + monkeypatch.setattr(board, "_commit_rekey", fail_commit) + with pytest.raises(RuntimeError, match="recovery"): + board.rekey_space("old", "new") + + assert tokens.space_rekeys() == {"old": "new"} + assert tokens.is_space_active("old") is False + with pytest.raises(ValueError, match="retired"): + tokens.mint("lead-1", "lead", space="old") + + board.close() + journal.close() + + +def test_existing_empty_case_gains_a_safe_space_association(board, journal): + journal.ensure_case("legacy-empty", WORKER.id) + + board.create_item( + "alpha", WORKER, title="Existing case", criteria="visible", case="legacy-empty" + ) + + assert journal.cases(WORKER, scope_space="alpha") == ["legacy-empty"] + + def test_explicit_grant_shares_across_teams(journal): journal.append(LEAD, "findings", "opened by lead-1") with pytest.raises(AuthorityError): diff --git a/tests/test_team_open_surface.py b/tests/test_team_open_surface.py index 6cde40d58..9ae1d210c 100644 --- a/tests/test_team_open_surface.py +++ b/tests/test_team_open_surface.py @@ -9,7 +9,7 @@ from coworker.teams import Actor, AuthorityError, BoardError, JournalStore, Role, TeamStore from coworker.teams.dialect import LocalDialect, RemoteDialect, local_dialect -from coworker.teams.tokens import BoardTokens +from coworker.teams.tokens import BoardPrincipal, BoardTokens USER = Actor(id="user", role=Role.USER) LEAD = Actor(id="lead-1", role=Role.LEAD, persona="swe-lead") @@ -200,14 +200,15 @@ def test_board_api_rejects_the_sidecar_token_as_a_board_token(api): def test_token_binds_identity_and_store_enforces_authority(api): client, manager, app = api - lead_token = _tokens(manager).mint("lead-1", "lead") - nia_token = _tokens(manager).mint("nia", "worker") + lead_token = _tokens(manager).mint("lead-1", "lead", space="proj") + nia_token = _tokens(manager).mint("nia", "worker", space="proj") lead = {"Authorization": f"Bearer {lead_token}"} nia = {"Authorization": f"Bearer {nia_token}"} assert client.get("/v1/board/whoami", headers=nia).json() == { "actor": "nia", "role": "worker", + "space": "proj", } created = client.post( @@ -249,11 +250,261 @@ def test_token_binds_identity_and_store_enforces_authority(api): assert bad.status_code == 403 or bad.status_code == 400 +def test_board_token_is_bound_to_one_space(api): + client, manager, _ = api + alpha_token = _tokens(manager).mint("lead-1", "lead", space="alpha") + alpha = {"Authorization": f"Bearer {alpha_token}"} + + assert client.get("/v1/board/whoami", headers=alpha).json() == { + "actor": "lead-1", + "role": "lead", + "space": "alpha", + } + assert client.get("/v1/board/spaces", headers=alpha).json() == { + "spaces": ["alpha"] + } + + for denied_space in ("beta", "does-not-exist"): + denied = client.get( + "/v1/board/items", headers=alpha, params={"space": denied_space} + ) + assert denied.status_code == 404 + assert denied.json() == {"error": "board space not found"} + + before = manager.team_store.event_count("beta") + denied_create = client.post( + "/v1/board/items", + headers=alpha, + json={"space": "beta", "title": "Denied", "criteria": "never written"}, + ) + assert denied_create.status_code == 404 + assert denied_create.json() == {"error": "board space not found"} + assert manager.team_store.event_count("beta") == before + + allowed = client.post( + "/v1/board/items", + headers=alpha, + json={"space": "alpha", "title": "Allowed", "criteria": "written"}, + ) + assert allowed.status_code == 200 + + +def test_wrong_space_is_rejected_before_every_board_handler(api): + client, manager, _ = api + token = _tokens(manager).mint("lead-1", "lead", space="alpha") + headers = {"Authorization": f"Bearer {token}"} + requests = [ + ("GET", "/v1/board/items", {"params": {"space": "beta"}}), + ("GET", "/v1/board/item", {"params": {"space": "beta", "id": 1}}), + ( + "GET", + "/v1/board/attachment", + {"params": {"space": "beta", "name": f"{'0' * 64}.png"}}, + ), + ("GET", "/v1/board/policy", {"params": {"space": "beta"}}), + ("GET", "/v1/board/pending", {"params": {"space": "beta"}}), + ( + "POST", + "/v1/board/items", + {"json": {"space": "beta", "title": "x", "criteria": "c"}}, + ), + ( + "POST", + "/v1/board/items/transition", + {"json": {"space": "beta", "id": 1, "to": "in_progress"}}, + ), + ( + "POST", + "/v1/board/items/comment", + {"json": {"space": "beta", "id": 1, "body": "x"}}, + ), + ( + "POST", + "/v1/board/items/assign", + {"json": {"space": "beta", "id": 1, "assignee": "nia"}}, + ), + ("POST", "/v1/board/items/claim", {"json": {"space": "beta", "id": 1}}), + ( + "POST", + "/v1/board/link", + {"json": {"space": "beta", "src": 1, "kind": "blocks", "dst": 2}}, + ), + ( + "POST", + "/v1/board/items/attach", + { + "json": { + "space": "beta", + "id": 1, + "filename": "x.png", + "data_b64": "not-base64", + } + }, + ), + ( + "POST", + "/v1/board/policy", + {"json": {"space": "beta", "claims": "lead-only"}}, + ), + ( + "POST", + "/v1/board/consume", + {"json": {"space": "beta", "upto_seq": 999}}, + ), + ( + "POST", + "/v1/board/journal", + {"json": {"space": "beta", "case": "secret", "body": "x"}}, + ), + ] + + before_events = manager.team_store.event_count("beta") + for method, path, kwargs in requests: + response = client.request(method, path, headers=headers, **kwargs) + assert response.status_code == 404, (method, path, response.text) + assert response.json() == {"error": "board space not found"} + assert manager.team_store.event_count("beta") == before_events + assert not manager.attachment_store.root.exists() + assert manager.journal_store.overview(USER) == [] + + +def test_board_token_scopes_journal_cases_reads_and_appends(api): + client, manager, _ = api + manager.journal_store.append(USER, "shared", "alpha finding", space="alpha") + manager.journal_store.append(USER, "shared", "beta finding", space="beta") + manager.journal_store.append(USER, "beta-only", "hidden", space="beta") + token = _tokens(manager).mint("user", "user", space="alpha") + headers = {"Authorization": f"Bearer {token}"} + + cases = client.get("/v1/board/journal/cases", headers=headers) + assert cases.status_code == 200 + assert [case["case"] for case in cases.json()["cases"]] == ["shared"] + assert cases.json()["cases"][0]["entries"] == 1 + + read = client.get( + "/v1/board/journal", headers=headers, params={"case": "shared"} + ) + assert read.status_code == 200 + assert [entry["body"] for entry in read.json()["entries"]] == [ + "alpha finding" + ] + + explicit_empty = client.post( + "/v1/board/journal", + headers=headers, + json={"space": "", "case": "shared", "body": "must be rejected"}, + ) + assert explicit_empty.status_code == 404 + assert explicit_empty.json() == {"error": "board space not found"} + + appended = client.post( + "/v1/board/journal", + headers=headers, + json={"case": "shared", "body": "scoped append"}, + ) + assert appended.status_code == 200 + assert appended.json()["space"] == "alpha" + assert [ + entry["body"] + for entry in manager.journal_store.read(USER, "shared", scope_space="alpha") + ] == ["alpha finding", "scoped append"] + + +def test_space_rekey_requires_a_new_token_and_moves_journal_projection(api): + client, manager, _ = api + old_token = _tokens(manager).mint("lead-1", "lead", space="old") + old = {"Authorization": f"Bearer {old_token}"} + item = client.post( + "/v1/board/items", + headers=old, + json={ + "space": "old", + "title": "Migrated", + "criteria": "still scoped", + "case": "migration", + }, + ).json() + assert client.post( + "/v1/board/journal", + headers=old, + json={"case": "migration", "body": "before rekey", "item": item["id"]}, + ).status_code == 200 + + assert manager.team_store.rekey_space("old", "new") is True + + assert client.get("/v1/board/whoami", headers=old).status_code == 401 + denied_fork = client.post( + "/v1/board/items", + headers=old, + json={"space": "old", "title": "Fork", "criteria": "must not exist"}, + ) + assert denied_fork.status_code == 401 + assert manager.team_store.event_count("old") == 0 + old_journal = client.get( + "/v1/board/journal", headers=old, params={"case": "migration"} + ) + assert old_journal.status_code == 401 + + new_token = _tokens(manager).mint("lead-1", "lead", space="new") + new = {"Authorization": f"Bearer {new_token}"} + assert client.get( + "/v1/board/item", + headers=new, + params={"space": "new", "id": item["id"]}, + ).status_code == 200 + moved_journal = client.get( + "/v1/board/journal", headers=new, params={"case": "migration"} + ) + assert [entry["body"] for entry in moved_journal.json()["entries"]] == [ + "before rekey" + ] + + +def test_space_rekey_rejects_a_principal_resolved_before_the_move(api, monkeypatch): + client, manager, _ = api + token = _tokens(manager).mint("lead-1", "lead", space="old") + principal = _tokens(manager).resolve(token) + manager.team_store.create_item("old", LEAD, title="Before", criteria="moves") + + assert principal is not None + assert manager.team_store.rekey_space("old", "new") is True + + # Simulate a request that resolved immediately before rekey and did not enter + # the store operation until immediately after it. + monkeypatch.setattr(manager.board_tokens, "resolve", lambda _: principal) + denied = client.post( + "/v1/board/items", + headers={"Authorization": f"Bearer {token}"}, + json={"space": "old", "title": "Fork", "criteria": "must not exist"}, + ) + + assert denied.status_code == 404 + assert denied.json() == {"error": "board space not found"} + assert manager.team_store.event_count("old") == 0 + with pytest.raises(ValueError, match="retired"): + manager.board_tokens.mint("lead-1", "lead", space="old") + + +def test_journal_append_rejects_explicit_null_space(api): + client, manager, _ = api + token = _tokens(manager).mint("user", "user", space="None") + + response = client.post( + "/v1/board/journal", + headers={"Authorization": f"Bearer {token}"}, + json={"space": None, "case": "null-scope", "body": "must not append"}, + ) + + assert response.status_code == 404 + assert response.json() == {"error": "board space not found"} + assert manager.journal_store.cases(USER) == [] + + def test_board_item_reads_hide_foreign_worker_items(api): client, manager, _ = api - lead_token = _tokens(manager).mint("lead-1", "lead") - nia_token = _tokens(manager).mint("nia", "worker") - webb_token = _tokens(manager).mint("webb", "worker") + lead_token = _tokens(manager).mint("lead-1", "lead", space="proj") + nia_token = _tokens(manager).mint("nia", "worker", space="proj") + webb_token = _tokens(manager).mint("webb", "worker", space="proj") lead = {"Authorization": f"Bearer {lead_token}"} nia = {"Authorization": f"Bearer {nia_token}"} webb = {"Authorization": f"Bearer {webb_token}"} @@ -356,9 +607,9 @@ def create(title, *, description=""): def test_board_item_reads_return_not_found_for_missing_items(api): client, manager, _ = api - nia_token = _tokens(manager).mint("nia", "worker") - lead_token = _tokens(manager).mint("lead-1", "lead") - user_token = _tokens(manager).mint("user", "user") + nia_token = _tokens(manager).mint("nia", "worker", space="proj") + lead_token = _tokens(manager).mint("lead-1", "lead", space="proj") + user_token = _tokens(manager).mint("user", "user", space="proj") nia = {"Authorization": f"Bearer {nia_token}"} lead = {"Authorization": f"Bearer {lead_token}"} user = {"Authorization": f"Bearer {user_token}"} @@ -374,8 +625,8 @@ def test_board_item_reads_return_not_found_for_missing_items(api): def test_board_item_create_hides_missing_and_foreign_parents(api): client, manager, _ = api - lead_token = _tokens(manager).mint("lead-1", "lead") - nia_token = _tokens(manager).mint("nia", "worker") + lead_token = _tokens(manager).mint("lead-1", "lead", space="proj") + nia_token = _tokens(manager).mint("nia", "worker", space="proj") lead = {"Authorization": f"Bearer {lead_token}"} nia = {"Authorization": f"Bearer {nia_token}"} @@ -411,8 +662,8 @@ def test_board_item_create_hides_missing_and_foreign_parents(api): def test_remote_dialect_round_trip(api): client, manager, app = api - lead_token = _tokens(manager).mint("lead-1", "lead") - nia_token = _tokens(manager).mint("nia", "worker") + lead_token = _tokens(manager).mint("lead-1", "lead", space="proj") + nia_token = _tokens(manager).mint("nia", "worker", space="proj") from fastapi.testclient import TestClient lead = RemoteDialect( @@ -465,8 +716,8 @@ def test_remote_dialect_round_trip(api): def test_pending_and_consume_over_the_wire(api): client, manager, app = api - lead_token = _tokens(manager).mint("lead-1", "lead") - nia_token = _tokens(manager).mint("nia", "worker") + lead_token = _tokens(manager).mint("lead-1", "lead", space="proj") + nia_token = _tokens(manager).mint("nia", "worker", space="proj") from fastapi.testclient import TestClient lead = RemoteDialect( @@ -525,8 +776,8 @@ def test_attach_over_the_wire_and_fetch(api): client, manager, app = api from fastapi.testclient import TestClient - lead_token = _tokens(manager).mint("lead-1", "lead") - nia_token = _tokens(manager).mint("nia", "worker") + lead_token = _tokens(manager).mint("lead-1", "lead", space="proj") + nia_token = _tokens(manager).mint("nia", "worker", space="proj") lead = RemoteDialect( "http://board.test", lead_token, @@ -556,7 +807,7 @@ def test_attach_over_the_wire_and_fetch(api): data, mime = lead.attachment("proj", stored) assert data == PNG and mime == "image/png" assert nia.attachment("proj", stored) == (PNG, "image/png") - with pytest.raises(BoardError, match="attachment not found"): + with pytest.raises(BoardError, match="board space not found"): lead.attachment("another-space", stored) # a worker cannot attach to an item outside its slice @@ -571,8 +822,8 @@ def test_board_attachment_read_hides_foreign_worker_reference(api): from fastapi.testclient import TestClient from coworker.teams.attachments import stored_name - lead_token = _tokens(manager).mint("lead-1", "lead") - nia_token = _tokens(manager).mint("nia", "worker") + lead_token = _tokens(manager).mint("lead-1", "lead", space="proj") + nia_token = _tokens(manager).mint("nia", "worker", space="proj") lead = {"Authorization": f"Bearer {lead_token}"} nia = {"Authorization": f"Bearer {nia_token}"} @@ -635,7 +886,7 @@ def test_board_attachment_read_requires_token(api): def test_board_attachment_read_rejects_malformed_name(api): client, manager, _ = api lead = { - "Authorization": f"Bearer {_tokens(manager).mint('lead-1', 'lead')}" + "Authorization": f"Bearer {_tokens(manager).mint('lead-1', 'lead', space='proj')}" } response = client.get( @@ -653,7 +904,7 @@ def test_board_attachment_read_hides_unreferenced_blob(api): from coworker.teams.attachments import stored_name lead = { - "Authorization": f"Bearer {_tokens(manager).mint('lead-1', 'lead')}" + "Authorization": f"Bearer {_tokens(manager).mint('lead-1', 'lead', space='proj')}" } ref = manager.attachment_store.put(PNG, "orphan.png") @@ -672,7 +923,7 @@ def test_board_attachment_read_hides_missing_referenced_blob(api): from coworker.teams.attachments import stored_name lead = { - "Authorization": f"Bearer {_tokens(manager).mint('lead-1', 'lead')}" + "Authorization": f"Bearer {_tokens(manager).mint('lead-1', 'lead', space='proj')}" } item = client.post( "/v1/board/items", @@ -862,7 +1113,7 @@ def test_local_attachment_read_allows_a_directly_linked_item(tmp_path): def test_attach_rejects_bad_payloads_over_the_wire(api): client, manager, app = api - lead_token = _tokens(manager).mint("lead-1", "lead") + lead_token = _tokens(manager).mint("lead-1", "lead", space="proj") headers = {"Authorization": f"Bearer {lead_token}"} client.post( "/v1/board/items", @@ -883,20 +1134,165 @@ def test_attach_rejects_bad_payloads_over_the_wire(api): def test_tokens_are_hash_stored_and_revocable(tmp_path): tokens = BoardTokens(tmp_path / "board-tokens.json") - token = tokens.mint("nia", "worker", label="laptop") + token = tokens.mint("nia", "worker", space="proj", label="laptop") # plaintext never touches disk assert token not in (tmp_path / "board-tokens.json").read_text() - actor = tokens.resolve(token) - assert (actor.id, actor.role) == ("nia", Role.WORKER) + principal = tokens.resolve(token) + assert principal == BoardPrincipal( + actor=Actor(id="nia", role=Role.WORKER), space="proj" + ) assert tokens.resolve("owb_forged") is None assert tokens.revoke(token[:12]) == 1 assert tokens.resolve(token) is None -def test_token_mint_validates_role(tmp_path): +@pytest.mark.parametrize("role", ["admin", "system"]) +def test_token_mint_validates_role(tmp_path, role): tokens = BoardTokens(tmp_path / "board-tokens.json") with pytest.raises(ValueError): - tokens.mint("nia", "admin") + tokens.mint("nia", role, space="proj") + + +@pytest.mark.parametrize("space", ["", " ", "*"]) +def test_token_mint_requires_one_exact_space(tmp_path, space): + tokens = BoardTokens(tmp_path / "board-tokens.json") + with pytest.raises(ValueError, match="space"): + tokens.mint("nia", "worker", space=space) + + +def test_token_space_is_an_opaque_exact_key(tmp_path): + tokens = BoardTokens(tmp_path / "board-tokens.json") + token = tokens.mint("nia", "worker", space="proj ") + + assert tokens.resolve(token).space == "proj " + + +def test_legacy_unscoped_tokens_fail_closed(tmp_path): + tokens = BoardTokens(tmp_path / "board-tokens.json") + token = tokens.mint("nia", "worker", space="proj") + path = tmp_path / "board-tokens.json" + entries = json.loads(path.read_text()) + next(iter(entries.values())).pop("space") + path.write_text(json.dumps(entries)) + + assert tokens.resolve(token) is None + + +@pytest.mark.parametrize( + ("field", "value"), + [ + ("actor", ""), + ("role", "admin"), + ("role", "system"), + ("space", "*"), + ("space", ["proj"]), + ], +) +def test_malformed_token_entries_fail_closed(tmp_path, field, value): + tokens = BoardTokens(tmp_path / "board-tokens.json") + token = tokens.mint("nia", "worker", space="proj") + path = tmp_path / "board-tokens.json" + entries = json.loads(path.read_text()) + next(iter(entries.values()))[field] = value + path.write_text(json.dumps(entries)) + + assert tokens.resolve(token) is None + + +def test_malformed_token_registry_containers_fail_closed(tmp_path): + tokens = BoardTokens(tmp_path / "board-tokens.json") + token = tokens.mint("nia", "worker", space="proj") + path = tmp_path / "board-tokens.json" + entries = json.loads(path.read_text()) + digest = next(iter(entries)) + + path.write_text("[]") + assert tokens.resolve(token) is None + + path.write_text(json.dumps({digest: []})) + assert tokens.resolve(token) is None + + +def test_malformed_rekey_tombstones_fail_closed(tmp_path): + tokens = BoardTokens(tmp_path / "board-tokens.json") + tokens.begin_space_rekey("old", "new") + path = tmp_path / "board-tokens.json" + entries = json.loads(path.read_text()) + entries["__space_rekeys__"] = [] + path.write_text(json.dumps(entries)) + + assert tokens.is_space_active("old") is False + with pytest.raises(ValueError, match="registry is malformed"): + tokens.mint("nia", "worker", space="old") + + +def test_token_rekey_tombstone_survives_a_cross_process_mint(tmp_path): + import subprocess + import sys + import threading + import time + + registry = tmp_path / "board-tokens.json" + loaded = tmp_path / "mint-loaded" + release = tmp_path / "release-mint" + token_output = tmp_path / "minted-token" + script = """ +import sys +import time +from pathlib import Path +from coworker.teams.tokens import BoardTokens + +registry, loaded, release, output = map(Path, sys.argv[1:]) +tokens = BoardTokens(registry) +original_load = tokens._load + +def paused_load(): + entries = original_load() + loaded.touch() + while not release.exists(): + time.sleep(0.01) + return entries + +tokens._load = paused_load +output.write_text(tokens.mint("nia", "worker", space="old")) +""" + process = subprocess.Popen( + [ + sys.executable, + "-c", + script, + str(registry), + str(loaded), + str(release), + str(token_output), + ], + stdout=subprocess.PIPE, + stderr=subprocess.PIPE, + text=True, + ) + deadline = time.monotonic() + 5 + while not loaded.exists() and process.poll() is None and time.monotonic() < deadline: + time.sleep(0.01) + assert loaded.exists(), process.communicate(timeout=1) + + tokens = BoardTokens(registry) + tombstoned = threading.Event() + + def begin_rekey(): + tokens.begin_space_rekey("old", "new") + tombstoned.set() + + thread = threading.Thread(target=begin_rekey) + thread.start() + assert tombstoned.wait(0.1) is False + release.touch() + stdout, stderr = process.communicate(timeout=5) + thread.join(timeout=5) + + assert process.returncode == 0, (stdout, stderr) + assert thread.is_alive() is False + assert tokens.space_rekeys() == {"old": "new"} + assert tokens.resolve(token_output.read_text()) is None # ------------------------------------------------------------------ MCP server @@ -927,6 +1323,23 @@ def names(server): assert "journal_append" in worker_names +def test_mcp_rejects_a_token_for_another_space(api): + from fastapi.testclient import TestClient + + from coworker.teams.mcp_server import build + + client, manager, app = api + token = _tokens(manager).mint("nia", "worker", space="alpha") + dialect = RemoteDialect( + "http://board.test", token, client=TestClient(app, base_url="http://board.test") + ) + try: + with pytest.raises(BoardError, match="token is scoped to 'alpha'"): + build(dialect, space="beta") + finally: + dialect.close() + + def test_mcp_worker_loop_through_call_tool(tmp_path): import anyio @@ -1106,10 +1519,80 @@ def test_cli_token_mint_and_list(tmp_path, capsys): assert main( ["board", "token", "mint", "--actor", "nia", "--role", "worker", - "--label", "laptop", "--db", str(tmp_path)] + "--space", "proj", "--label", "laptop", "--db", str(tmp_path)] ) == 0 token = capsys.readouterr().out.strip() assert token.startswith("owb_") assert main(["board", "token", "list", "--db", str(tmp_path)]) == 0 out = capsys.readouterr().out - assert "nia" in out and "laptop" in out and token not in out + assert "nia" in out and "proj" in out and "laptop" in out and token not in out + + +def test_cli_token_mint_requires_space(tmp_path, capsys): + from coworker.teams.cli import main + + result = main( + ["board", "token", "mint", "--actor", "nia", "--db", str(tmp_path)] + ) + + assert result == 1 + assert "--space is required" in capsys.readouterr().err + + +def test_cli_token_mint_rejects_wildcard_space(tmp_path, capsys): + from coworker.teams.cli import main + + result = main( + [ + "board", + "token", + "mint", + "--actor", + "nia", + "--space", + "*", + "--db", + str(tmp_path), + ] + ) + + assert result == 1 + assert "exact board space" in capsys.readouterr().err + + +def test_cli_token_mint_reports_a_retired_space_cleanly(tmp_path, capsys): + from coworker.teams.cli import main + + tokens = BoardTokens(tmp_path / "board-tokens.json") + tokens.begin_space_rekey("old", "new") + + result = main( + [ + "board", + "token", + "mint", + "--actor", + "nia", + "--space", + "old", + "--db", + str(tmp_path), + ] + ) + + assert result == 1 + assert "space 'old' is retired" in capsys.readouterr().err + + +def test_local_cli_caches_a_distinct_token_per_space(): + from coworker.secrets import state_dir + from coworker.teams.cli import _local_cli_token + + alpha = _local_cli_token("alpha") + assert _local_cli_token("alpha") == alpha + beta = _local_cli_token("beta") + assert beta != alpha + + tokens = BoardTokens(state_dir() / "board-tokens.json") + assert tokens.resolve(alpha).space == "alpha" + assert tokens.resolve(beta).space == "beta"