Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
2 changes: 1 addition & 1 deletion Dockerfile
Original file line number Diff line number Diff line change
Expand Up @@ -4,7 +4,7 @@ FROM node:22-bookworm-slim@sha256:6c74791e557ce11fc957704f6d4fe134a7bc8d6f5ca440

WORKDIR /build/sdk/typescript

COPY sdk/typescript/package.json sdk/typescript/pnpm-lock.yaml ./
COPY sdk/typescript/package.json sdk/typescript/pnpm-lock.yaml sdk/typescript/pnpm-workspace.yaml ./
COPY plugins/codex-security/mcp-app/package.json plugins/codex-security/mcp-app/package-lock.json /build/plugins/codex-security/mcp-app/

RUN corepack enable \
Expand Down
6 changes: 4 additions & 2 deletions README.md
Original file line number Diff line number Diff line change
Expand Up @@ -50,8 +50,10 @@ Use the included Docker Compose configuration for scans of many repositories. Se

The [findings service](sdk/typescript/README.md#findings-service-preview) runs
from the SDK in Docker, stores findings and embeddings in SQLite, and lists
findings with pagination. It also returns potential duplicates by embedding
similarity within a repository or an explicit all-repository scope. The
findings with pagination. Its read-only dashboard at `/dashboard` refreshes every
five seconds and shows stored findings and duplicate groups from the service's
database. It also returns potential duplicates by embedding similarity within a
repository or an explicit all-repository scope. The
`codex-security publish scan --to custom --findings-url http://localhost:3000`
command uploads completed findings and their repository ID. The SDK and
`codex-security dedupe` command retrieve candidates, run independent Codex
Expand Down
1 change: 1 addition & 0 deletions plugins/codex-security/scripts/workbench_cli.py
Original file line number Diff line number Diff line change
Expand Up @@ -339,6 +339,7 @@ def parse_args(description: str) -> argparse.Namespace:
publication.add_argument("--input-file", required=True)

subparsers.add_parser("database-info")
subparsers.add_parser("dashboard")
subparsers.add_parser("finding-workflow")
subparsers.add_parser("store-findings")
subparsers.add_parser("store-dedupe-groups")
Expand Down
115 changes: 115 additions & 0 deletions plugins/codex-security/scripts/workbench_dashboard.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,115 @@
"""Read-only dashboard projections for stored findings and duplicate groups."""

from __future__ import annotations

import argparse
import json
import sqlite3
import sys
from pathlib import Path
from typing import Any

sys.path.insert(0, str(Path(__file__).resolve().parent))
from workbench_findings import list_dedupe_groups


FINDING_RECORDS = """
SELECT findings.id, json_extract(details_json, '$.title') AS title,
COALESCE(repositories.ids, '[]') AS repositoryIds,
json_extract(details_json, '$.severity.level') AS severity,
findings.created_at AS createdAt, findings.updated_at AS updatedAt
FROM findings LEFT JOIN (
SELECT finding_id, json_group_array(repository_id) AS ids
FROM finding_repositories GROUP BY finding_id
) AS repositories ON repositories.finding_id = findings.id
WHERE details_json IS NOT NULL
"""

GROUP_RECORDS = """
SELECT groups.id, groups.id AS title,
(SELECT json_group_array(DISTINCT repository_id)
FROM finding_dedupe_group_members AS members
JOIN finding_repositories ON finding_repositories.finding_id = members.finding_id
WHERE members.group_id = groups.id) AS repositoryIds,
groups.created_at AS createdAt, groups.created_at AS updatedAt,
(SELECT COUNT(*) FROM finding_dedupe_group_members WHERE group_id = groups.id) AS memberCount
FROM finding_dedupe_groups AS groups
"""

RECORDS = {
"findings": FINDING_RECORDS,
"groups": GROUP_RECORDS,
}


def item(row: sqlite3.Row) -> dict[str, Any]:
result = dict(row)
result["repositoryIds"] = json.loads(result["repositoryIds"])
return result


def detail(connection: sqlite3.Connection, view: str, selected: dict[str, Any]) -> dict[str, Any]:
result: dict[str, Any] = {"item": selected}
selected_id = selected["id"]
if view == "findings":
result["finding"] = json.loads(connection.execute(
"SELECT details_json FROM findings WHERE id = ?", (selected_id,),
).fetchone()[0])
result["groups"] = list_dedupe_groups(connection, selected_id)["groups"]
else:
result["group"] = {
"groupId": selected_id, "createdAt": selected["createdAt"],
"findingIds": [r[0] for r in connection.execute(
"SELECT finding_id FROM finding_dedupe_group_members WHERE group_id = ? ORDER BY finding_id",
(selected_id,),
)],
}
return result


def dashboard(connection: sqlite3.Connection, query: dict[str, Any]) -> dict[str, Any]:
"""One snapshot, no artifact reads, model calls, or writes."""
view = query["view"]
records = RECORDS[view]
clauses: list[str] = []
values: list[Any] = []
if query.get("query"):
connection.create_function("casefold", 1, str.casefold, deterministic=True)
columns = ["id", "title", "repositoryIds"]
clauses.append("(" + " OR ".join(f"instr(casefold(COALESCE({c}, '')), casefold(?)) > 0" for c in columns) + ")")
values.extend([query["query"]] * len(columns))
if query.get("repository"):
clauses.append("EXISTS (SELECT 1 FROM json_each(repositoryIds) WHERE value = ?)")
values.append(query["repository"])
where = " WHERE " + " AND ".join(clauses) if clauses else ""
order = "createdAt DESC, id" if query["sort"] == "newest" else "updatedAt DESC, id"
connection.execute("BEGIN")
with connection:
repositories = connection.execute("""
SELECT DISTINCT repository_id AS id, repository_id AS label
FROM finding_repositories ORDER BY repository_id
""").fetchall()
total = connection.execute(f"SELECT COUNT(*) FROM ({records}) {where}", values).fetchone()[0]
rows = connection.execute(
f"SELECT * FROM ({records}) {where} ORDER BY {order} LIMIT ? OFFSET ?",
(*values, query["limit"], query["offset"]),
).fetchall()
selected = connection.execute(
f"SELECT * FROM ({records}) WHERE id = ?", (query["id"],),
).fetchone() if query.get("id") else None
next_offset = query["offset"] + len(rows)
return {
"overview": {
"findings": connection.execute("SELECT COUNT(*) FROM findings WHERE details_json IS NOT NULL").fetchone()[0],
"groups": connection.execute("SELECT COUNT(*) FROM finding_dedupe_groups").fetchone()[0],
},
"repositories": [dict(row) for row in repositories],
"items": [item(row) for row in rows], "total": total,
"limit": query["limit"], "offset": query["offset"],
"nextOffset": next_offset if next_offset < total else None,
"detail": detail(connection, view, item(selected)) if selected is not None else None,
}


if __name__ == "__main__":
argparse.ArgumentParser(description=__doc__).parse_args()
3 changes: 3 additions & 0 deletions plugins/codex-security/scripts/workbench_db.py
Original file line number Diff line number Diff line change
Expand Up @@ -81,6 +81,7 @@
SQLITE_RETRY_ATTEMPTS,
)
from workbench_feedback import get_scan_feedback
from workbench_dashboard import dashboard
from workbench_finding_index import index_findings
from workbench_finding_workflows import finding_workflow, register_workflow_scan
from workbench_findings import (
Expand Down Expand Up @@ -4038,6 +4039,8 @@ def main() -> None:
result = {"databasePath": str(database_path())}
elif args.command == "finding-workflow":
result = finding_workflow(connection, json.load(sys.stdin), now())
elif args.command == "dashboard":
result = dashboard(connection, json.load(sys.stdin))
elif args.command == "store-findings":
payload = json.load(sys.stdin)
result = store_findings(connection, payload["entries"], now(), payload.get("repositoryId"))
Expand Down
44 changes: 44 additions & 0 deletions plugins/codex-security/scripts/workbench_finding_workflows.py
Original file line number Diff line number Diff line change
Expand Up @@ -4,9 +4,15 @@

import argparse
import json
import hashlib
import sqlite3
import sys
from pathlib import Path
from typing import Any

sys.path.insert(0, str(Path(__file__).resolve().parent))
from workbench_target import directory_content_digest, git_output, git_revision

WORKFLOW_BINDINGS = {
"repositoryPath": "repository_path",
"scanRequestDigest": "scan_request_digest",
Expand Down Expand Up @@ -114,6 +120,42 @@ def finding_workflow(
raise SystemExit("workflowId must be a nonempty string.")
if payload["action"] == "get":
return {"workflow": read_workflow(connection, workflow_id)}
if payload["action"] == "source":
Comment thread
kmbroai marked this conversation as resolved.
target = Path(payload["repository"]).resolve(strict=True)
return {"source": {
"repository": str(target),
"revision": git_revision(target),
"refsDigest": hashlib.sha256((git_output(target, "show-ref") or "").encode()).hexdigest(),
"content": directory_content_digest(target, include_ignored=True),
}}
if payload["action"] == "get-review":
row = connection.execute(
"SELECT result_json FROM finding_workflow_reviews WHERE workflow_id = ? AND review_key = ?",
(workflow_id, payload["key"]),
).fetchone()
return {"review": json.loads(row["result_json"]) if row is not None else None}
if payload["action"] == "save-review":
binding = payload["binding"]
source = binding["source"]
scope = binding["scope"]
with connection:
connection.execute(
"""INSERT INTO finding_workflow_reviews
(workflow_id, review_key, review_contract_version, codex_version,
source_repository_path, source_revision, source_refs_digest, source_content_digest,
scope_repository_id, scope_all_repositories, model, effort, settings_digest,
prompt_digest, contract_digest, result_json, created_at)
VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?)
ON CONFLICT(workflow_id, review_key) DO NOTHING""",
(
workflow_id, payload["key"], binding["version"], binding["codexVersion"],
source["repository"], source["revision"], source["refsDigest"], source["content"],
scope.get("repositoryId"), scope.get("allRepositories"), binding["model"],
binding["effort"], binding.get("settingsDigest"), binding["promptDigest"],
binding["contractDigest"], json.dumps(payload["result"], allow_nan=False), timestamp,
),
)
return {}
connection.execute("BEGIN IMMEDIATE")
with connection:
state = read_workflow(connection, workflow_id)
Expand All @@ -136,6 +178,8 @@ def finding_workflow(
state["stages"][stage] = {"status": "completed", "result": payload["result"]}
elif action == "fail":
current.update(status="failed", error=payload["error"])
elif action == "prepare-dedupe":
current.update(result=payload["result"], pendingWrite=payload["pendingWrite"])
else:
raise SystemExit("Unknown workflow action.")
save_workflow(connection, state, timestamp)
Expand Down
59 changes: 58 additions & 1 deletion plugins/codex-security/scripts/workbench_schema.py
Original file line number Diff line number Diff line change
Expand Up @@ -771,7 +771,20 @@
);
""",
),
# Version 37 is the stacked dedupe-review checkpoint migration.
(
37,
Comment thread
kmbroai marked this conversation as resolved.
"checkpoint validated dedupe reviews",
"""
CREATE TABLE finding_workflow_reviews (
workflow_id TEXT NOT NULL REFERENCES finding_workflows(id) ON DELETE CASCADE,
review_key TEXT NOT NULL,
binding_json TEXT NOT NULL,
result_json TEXT NOT NULL,
created_at TEXT NOT NULL,
PRIMARY KEY (workflow_id, review_key)
);
""",
),
(
38,
"store findings workflow metadata in columns",
Expand All @@ -793,9 +806,51 @@
ALTER TABLE finding_workflows ADD COLUMN dedupe_error TEXT;
""",
),
(
39,
"store dedupe checkpoint bindings in columns",
"""
ALTER TABLE finding_workflow_reviews RENAME COLUMN binding_json TO prompt_digest;
ALTER TABLE finding_workflow_reviews ADD COLUMN review_contract_version INTEGER;
ALTER TABLE finding_workflow_reviews ADD COLUMN codex_version TEXT;
ALTER TABLE finding_workflow_reviews ADD COLUMN source_repository_path TEXT;
ALTER TABLE finding_workflow_reviews ADD COLUMN source_revision TEXT;
ALTER TABLE finding_workflow_reviews ADD COLUMN source_refs_digest TEXT;
ALTER TABLE finding_workflow_reviews ADD COLUMN source_content_digest TEXT;
ALTER TABLE finding_workflow_reviews ADD COLUMN scope_repository_id TEXT;
ALTER TABLE finding_workflow_reviews ADD COLUMN scope_all_repositories INTEGER;
ALTER TABLE finding_workflow_reviews ADD COLUMN model TEXT;
ALTER TABLE finding_workflow_reviews ADD COLUMN effort TEXT;
ALTER TABLE finding_workflow_reviews ADD COLUMN settings_digest TEXT;
ALTER TABLE finding_workflow_reviews ADD COLUMN contract_digest TEXT;
""",
),
)


def migrate_finding_workflow_review_columns(connection: sqlite3.Connection) -> None:
for row in connection.execute(
"SELECT workflow_id, review_key, prompt_digest FROM finding_workflow_reviews"
).fetchall():
binding = json.loads(row["prompt_digest"])
source = binding["source"]
scope = binding["scope"]
connection.execute(
"""UPDATE finding_workflow_reviews SET review_contract_version = ?, codex_version = ?,
source_repository_path = ?, source_revision = ?, source_refs_digest = ?,
source_content_digest = ?, scope_repository_id = ?, scope_all_repositories = ?,
model = ?, effort = ?, settings_digest = ?, prompt_digest = ?, contract_digest = ?
WHERE workflow_id = ? AND review_key = ?""",
(
binding["version"], binding["codexVersion"], source["repository"],
source["revision"], source["refsDigest"], source["content"],
scope.get("repositoryId"), scope.get("allRepositories"), binding["model"],
binding["effort"], binding.get("settingsDigest"), binding["promptDigest"],
binding["contractDigest"], row["workflow_id"], row["review_key"],
),
)


def migrate_finding_workflow_columns(connection: sqlite3.Connection) -> None:
# Rename/backfill in place so existing checkpoint foreign keys and rows survive.
for row in connection.execute("SELECT id, results_json FROM finding_workflows").fetchall():
Expand Down Expand Up @@ -905,6 +960,8 @@ def apply_migrations(
connection.execute(statement)
if version == 38:
migrate_finding_workflow_columns(connection)
elif version == 39:
migrate_finding_workflow_review_columns(connection)
connection.execute(
"INSERT INTO schema_migrations (version, name, applied_at) VALUES (?, ?, ?)",
(version, name, now()),
Expand Down
25 changes: 22 additions & 3 deletions plugins/codex-security/scripts/workbench_target.py
Original file line number Diff line number Diff line change
Expand Up @@ -414,14 +414,31 @@ def git_directory_snapshot_paths(target: Path) -> list[Path] | None:
return sorted(set(paths))


def directory_content_digest(target: Path, *, excluded: tuple[Path, ...] = ()) -> str:
def source_directory_snapshot_paths(target: Path) -> list[Path]:
paths: list[Path] = []
pending = [target]
while pending:
for path in pending.pop().iterdir():
Comment thread
kmbroai marked this conversation as resolved.
if path.name == ".git":
continue
paths.append(path)
metadata = path.lstat()
# Name-surrogate reparse points include Windows directory junctions.
if stat.S_ISDIR(metadata.st_mode) and not getattr(metadata, "st_reparse_tag", 0) & 0x20000000:
pending.append(path)
return sorted(paths)


def directory_content_digest(
target: Path, *, excluded: tuple[Path, ...] = (), include_ignored: bool = False
) -> str:
excluded_relative = []
for path in excluded:
try:
excluded_relative.append(path.relative_to(target))
except ValueError:
continue
paths = git_directory_snapshot_paths(target)
paths = source_directory_snapshot_paths(target) if include_ignored else git_directory_snapshot_paths(target)
if paths is None:
paths = sorted(target.rglob("*"))
digest = hashlib.sha256()
Expand All @@ -440,7 +457,9 @@ def directory_content_digest(target: Path, *, excluded: tuple[Path, ...] = ()) -
raw_path = os.fsencode(relative_path.as_posix())
update_digest_field(digest, b"path", raw_path)
update_digest_field(digest, b"mode", str(stat.S_IMODE(metadata.st_mode)).encode())
if stat.S_ISLNK(metadata.st_mode):
if stat.S_ISLNK(metadata.st_mode) or (
include_ignored and getattr(metadata, "st_reparse_tag", 0) & 0x20000000
):
update_digest_field(digest, b"kind", b"symlink")
update_digest_field(digest, b"content", os.fsencode(os.readlink(path)))
elif stat.S_ISDIR(metadata.st_mode):
Expand Down
3 changes: 2 additions & 1 deletion plugins/codex-security/tests/test_workbench_db.py
Original file line number Diff line number Diff line change
Expand Up @@ -67,6 +67,7 @@
"finding_remediation_attempts",
"finding_repositories",
"finding_triage",
"finding_workflow_reviews",
"finding_workflows",
"findings",
"scan_artifacts",
Expand Down Expand Up @@ -995,7 +996,7 @@ def test_workbench_persists_progress_and_indexes_completed_findings(tmp_path: Pa
)
}
assert tables == EXPECTED_TABLES
assert connection.execute("SELECT COUNT(*) FROM schema_migrations").fetchone() == (37,)
assert connection.execute("SELECT COUNT(*) FROM schema_migrations").fetchone() == (39,)
assert connection.execute("SELECT COUNT(*) FROM findings").fetchone() == (1,)
assert connection.execute("SELECT COUNT(*) FROM finding_locations").fetchone() == (1,)

Expand Down
2 changes: 1 addition & 1 deletion plugins/codex-security/tests/test_workbench_deep_scan.py
Original file line number Diff line number Diff line change
Expand Up @@ -278,7 +278,7 @@ def claim() -> dict[str, object]:
return claim_deep_scan_coordinator(state_dir, codex_home, scan_id)

with sqlite3.connect(state_dir / "workbench.sqlite3") as connection:
assert connection.execute("SELECT MAX(version) FROM schema_migrations").fetchone() == (38,)
assert connection.execute("SELECT MAX(version) FROM schema_migrations").fetchone() == (39,)
assert claim()["deepScan"]["coordinatorGeneration"] == 2
assert claim()["coordinatorDisposition"] == "observing"
expire_deep_scan_coordinator(state_dir, scan_id)
Expand Down
Loading
Loading