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 application/single_app/config.py
Original file line number Diff line number Diff line change
Expand Up @@ -101,7 +101,7 @@
EXECUTOR_TYPE = 'thread'
EXECUTOR_MAX_WORKERS = 30
SESSION_TYPE = 'filesystem'
VERSION = "0.261.295"
VERSION = "0.261.296"
IS_DEVELOPMENT = is_development_env_enabled()

# Opt-out for deployments where App Service Easy Auth is active but the platform
Expand Down
119 changes: 95 additions & 24 deletions application/single_app/functions_control_center_activity.py
Original file line number Diff line number Diff line change
Expand Up @@ -14,6 +14,11 @@
ACTIVITY_PAGE_MAX = 200
ACTIVITY_SUMMARY_MAX = 5000
ACTIVITY_EXPORT_MAX = 10000
ACTIVITY_EXPORT_HEADER = (
"timestamp", "id", "user_id", "activity_type", "workspace_type", "user_name", "user_email",
"activity", "summary", "workspace_id", "workspace_name", "raw_json",

Check warning on line 19 in application/single_app/functions_control_center_activity.py

View workflow job for this annotation

GitHub Actions / malicious-pr-security-review

Important - Changed line contains AI, plugin, agent, or workspace boundary marker. Recommendation%3A Check whether prompts, chat history, uploaded documents, embeddings, citations, settings, or identity can cross a new boundary.
)
ACTIVITY_EXPORT_READABLE_COLUMNS = 6
ACTIVITY_ORDER = " ORDER BY c.timestamp DESC, c.id DESC, c.user_id DESC"
ACTIVITY_COMPOSITE_INDEX = [
{"path": "/timestamp", "order": "descending"},
Expand All @@ -28,6 +33,44 @@
"group.group_name", "workspace_context.group_id",
"workspace_context.public_workspace_id",
)
# Cosmos SQL keywords that cannot be used as dotted property names; the query grammar only

Check warning on line 36 in application/single_app/functions_control_center_activity.py

View workflow job for this annotation

GitHub Actions / malicious-pr-security-review

Important - Changed line contains secret or sensitive data source marker. Recommendation%3A Pair this source with any nearby network, logging, serialization, or process execution sink before approving.
# accepts ALL, FIRST and LAST there. A dotted reserved word (c.group.group_id) is a syntax
# error that rejects the whole query, so those segments are written as c['group'].
COSMOS_RESERVED_WORDS = frozenset({

Check warning on line 39 in application/single_app/functions_control_center_activity.py

View workflow job for this annotation

GitHub Actions / malicious-pr-security-review

Important - Changed line contains secret or sensitive data source marker. Recommendation%3A Pair this source with any nearby network, logging, serialization, or process execution sink before approving.
"AND", "ARRAY", "AS", "ASC", "BETWEEN", "BY", "DESC", "DISTINCT", "ESCAPE", "EXISTS",
"FALSE", "FROM", "GROUP", "IN", "JOIN", "LEFT", "LIKE", "LIMIT", "NOT", "NULL", "OFFSET",
"OR", "ORDER", "RANK", "RIGHT", "SELECT", "TOP", "TRUE", "UDF", "UNDEFINED", "VALUE", "WHERE",
})


def cosmos_property_path(path, root="c"):

Check warning on line 46 in application/single_app/functions_control_center_activity.py

View workflow job for this annotation

GitHub Actions / malicious-pr-security-review

Important - Changed line contains secret or sensitive data source marker. Recommendation%3A Pair this source with any nearby network, logging, serialization, or process execution sink before approving.
"""Return a Cosmos property reference that is valid even when a segment is a keyword."""

Check warning on line 47 in application/single_app/functions_control_center_activity.py

View workflow job for this annotation

GitHub Actions / malicious-pr-security-review

Important - Changed line contains secret or sensitive data source marker. Recommendation%3A Pair this source with any nearby network, logging, serialization, or process execution sink before approving.
expression = root
for segment in path.split("."):
expression += f"['{segment}']" if segment.upper() in COSMOS_RESERVED_WORDS else f".{segment}"

Check warning on line 50 in application/single_app/functions_control_center_activity.py

View workflow job for this annotation

GitHub Actions / malicious-pr-security-review

Important - Changed line contains secret or sensitive data source marker. Recommendation%3A Pair this source with any nearby network, logging, serialization, or process execution sink before approving.
return expression


NESTED_GROUP_ID = cosmos_property_path("group.group_id")

Check warning on line 54 in application/single_app/functions_control_center_activity.py

View workflow job for this annotation

GitHub Actions / malicious-pr-security-review

Important - Changed line contains secret or sensitive data source marker. Recommendation%3A Pair this source with any nearby network, logging, serialization, or process execution sink before approving.
# Where the writers record who acted. Records from approvals, membership and status changes
# often have no top-level user_id, so the person filter and the people search match all of

Check warning on line 56 in application/single_app/functions_control_center_activity.py

View workflow job for this annotation

GitHub Actions / malicious-pr-security-review

Important - Changed line contains secret or sensitive data source marker. Recommendation%3A Pair this source with any nearby network, logging, serialization, or process execution sink before approving.
# them. The display module resolves the table's Person column from the same fields.
ACTIVITY_ACTOR_FIELDS = (
"user_id", "admin_user_id", "requester_id", "added_by_user_id", "changed_by_user_id",
"changed_by.user_id", "removed_by.user_id", "admin.user_id", "actor.user_id",
)
# Where the writers record a group or public workspace. Status changes and member removals
# nest it (group.group_id, public_workspace.*), approvals and membership audits store it at

Check warning on line 63 in application/single_app/functions_control_center_activity.py

View workflow job for this annotation

GitHub Actions / malicious-pr-security-review

Important - Changed line contains AI, plugin, agent, or workspace boundary marker. Recommendation%3A Check whether prompts, chat history, uploaded documents, embeddings, citations, settings, or identity can cross a new boundary.

Check warning on line 63 in application/single_app/functions_control_center_activity.py

View workflow job for this annotation

GitHub Actions / malicious-pr-security-review

Important - Changed line contains security control, sanitization, or audit marker. Recommendation%3A Confirm the change does not weaken auth, CSRF, CSP, XSS defenses, settings sanitization, redaction, audit logging, or tests.
# the top level, user agreements store workspace_context.<type>_workspace_id, and public
# workspace ownership approvals store only a bare workspace_id. The display module resolves
# the Workspace column from the same locations.
GROUP_REFERENCE_FIELDS = (
"workspace_context.group_id", "group_id", "group.group_id", "workspace_context.group_workspace_id",
)
PUBLIC_REFERENCE_FIELDS = (
"workspace_context.public_workspace_id", "public_workspace_id", "public_workspace.public_workspace_id",
"public_workspace.workspace_id", "workspace_id",
)


def parse_activity_filters(args, now=None):
Expand Down Expand Up @@ -60,7 +103,14 @@
return result


def activity_query_context(filters):
def _either(fields, template):
"""OR the same predicate over several property paths, bracket-quoting reserved words."""
return " OR ".join(template.format(path=cosmos_property_path(field)) for field in fields)


def activity_query_context(filters, search_user_ids=()):
"""Build the WHERE clause. search_user_ids are people whose name or email matched the
search; they widen the search only and stay outside the cursor's filter scope."""
end_exclusive = (date.fromisoformat(filters["end_date"]) + timedelta(days=1)).isoformat()
clauses = ["IS_STRING(c.timestamp)", "IS_STRING(c.id)",
"(IS_STRING(c.user_id) OR IS_NULL(c.user_id) OR NOT IS_DEFINED(c.user_id))",
Expand All @@ -74,21 +124,22 @@
parameters.append({"name": name, "value": value})

add("ARRAY_CONTAINS(@types, c.activity_type)", "@types", filters["activity_types"])
add("(c.user_id = @user OR c.changed_by.user_id = @user OR c.admin_user_id = @user)",
"@user", filters["user_id"])
add(f"({_either(ACTIVITY_ACTOR_FIELDS, '{path} = @user')})", "@user", filters["user_id"])
workspace_type = filters["workspace_type"]
if workspace_type == "public":
add("c.workspace_type IN ('public', 'public_workspace')", "@workspace_type", workspace_type)
else:
add("c.workspace_type = @workspace_type", "@workspace_type", workspace_type)
group_id = filters["group_id"] or (filters["workspace_id"] if workspace_type == "group" else "")
public_id = filters["public_workspace_id"] or (filters["workspace_id"] if workspace_type == "public" else "")
add("(c.workspace_context.group_id = @group OR c.group_id = @group OR c.group.group_id = @group)",
"@group", group_id)
add("(c.workspace_context.public_workspace_id = @public OR c.public_workspace_id = @public)",
"@public", public_id)
# A specific group or public workspace matches every record that references it, however
# its writer recorded the workspace type; a type on its own matches the same references.
if workspace_type == "personal":
add("c.workspace_type = @workspace_type", "@workspace_type", workspace_type)
add("c.user_id = @personal", "@personal", filters["workspace_id"])
elif workspace_type == "group" and not group_id:
clauses.append(f"(c.workspace_type = 'group' OR {_either(GROUP_REFERENCE_FIELDS, 'IS_STRING({path})')})")
elif workspace_type == "public" and not public_id:
clauses.append("(c.workspace_type IN ('public', 'public_workspace') OR "
f"{_either(PUBLIC_REFERENCE_FIELDS, 'IS_STRING({path})')})")
add(f"({_either(GROUP_REFERENCE_FIELDS, '{path} = @group')})", "@group", group_id)
add(f"({_either(PUBLIC_REFERENCE_FIELDS, '{path} = @public')})", "@public", public_id)
add("c.token_type = @token_type", "@token_type", filters["token_type"])
add("c.usage.model = @model", "@model", filters["model"])
if filters["status"] == "failed":
Expand All @@ -99,7 +150,10 @@
add("(c.status = @status OR c.status_change.new_status = @status OR c.document.status = @status)",
"@status", filters["status"])
if filters["search"]:
search = " OR ".join(f"CONTAINS(c.{field}, @search, true)" for field in ACTIVITY_SEARCH_FIELDS)
search = _either(ACTIVITY_SEARCH_FIELDS, "CONTAINS({path}, @search, true)")
if search_user_ids:
search += " OR " + _either(ACTIVITY_ACTOR_FIELDS, "ARRAY_CONTAINS(@search_people, {path})")
parameters.append({"name": "@search_people", "value": list(search_user_ids)})
add(f"({search})", "@search", filters["search"])
return " AND ".join(clauses), parameters

Expand Down Expand Up @@ -138,9 +192,9 @@
raise ValueError("Invalid activity cursor") from ex


def query_activity_rows(container, filters, limit, cursor=None, snapshot=None, projection="*"):
def query_activity_rows(container, filters, limit, cursor=None, snapshot=None, projection="*", search_user_ids=()):
"""Three-part key: IDs are only unique inside a user partition, not globally."""
where, parameters = activity_query_context(filters)
where, parameters = activity_query_context(filters, search_user_ids)
snapshot = cursor["snapshot"] if cursor else snapshot or datetime.now(timezone.utc).isoformat()
# Compare calendar instants from both legacy naive-UTC and aware-UTC writers.
# The cursor keeps the stored timestamp verbatim so lexical ordering and seeking agree.
Expand Down Expand Up @@ -170,11 +224,13 @@
return rows, snapshot


def activity_page(container, filters, page_size=50, cursor_value=None):
def activity_page(container, filters, page_size=50, cursor_value=None, search_user_ids=()):
if not 1 <= page_size <= ACTIVITY_PAGE_MAX:
raise ValueError("Invalid page size")
cursor = decode_activity_cursor(cursor_value, filters) if cursor_value else None
rows, snapshot = query_activity_rows(container, filters, page_size + 1, cursor=cursor)
rows, snapshot = query_activity_rows(
container, filters, page_size + 1, cursor=cursor, search_user_ids=search_user_ids,
)
more = len(rows) > page_size
records = rows[:page_size]
return {
Expand All @@ -183,10 +239,11 @@
}


def activity_summary(container, filters):
def activity_summary(container, filters, search_user_ids=()):
"""Bounded projection, not COUNT/GROUP BY scans. Disclose sampling in the contract."""
rows, snapshot = query_activity_rows(
container, filters, ACTIVITY_SUMMARY_MAX + 1, projection="c.timestamp, c.id, c.user_id, c.activity_type",
search_user_ids=search_user_ids,
)
truncated = len(rows) > ACTIVITY_SUMMARY_MAX
rows = rows[:ACTIVITY_SUMMARY_MAX]
Expand Down Expand Up @@ -216,27 +273,41 @@
return text


def activity_csv_stream(container, filters, first_rows, snapshot):
"""Stream page-sized reads with a final status row when the export hits its cap."""
def activity_csv_stream(container, filters, first_rows, snapshot, search_user_ids=(), present_rows=None):
"""Stream page-sized reads with a final status row when the export hits its cap.

The first five columns and the trailing raw JSON keep their 0.261.284 positions. The
readable columns between them come from present_rows(rows), which returns one
(user_name, user_email, activity, summary, workspace_id, workspace_name) per row.
"""
output = StringIO()
writer = csv.writer(output)
blank = ("",) * ACTIVITY_EXPORT_READABLE_COLUMNS

def line(values):
output.seek(0)
output.truncate(0)
writer.writerow([activity_csv_cell(value) for value in values])
return output.getvalue()

yield line(("timestamp", "id", "user_id", "activity_type", "workspace_type", "raw_json"))
yield line(ACTIVITY_EXPORT_HEADER)
rows = first_rows
count = 0
while rows:
for row in rows[:ACTIVITY_EXPORT_MAX - count]:
batch = rows[:ACTIVITY_EXPORT_MAX - count]
readable = list(present_rows(batch)) if present_rows else []
if len(readable) != len(batch):
readable = [blank] * len(batch)
for row, columns in zip(batch, readable):
yield line((row.get("timestamp"), row.get("id"), row.get("user_id"), row.get("activity_type"),
row.get("workspace_type"), json.dumps(row, ensure_ascii=False)))
row.get("workspace_type"), *columns, json.dumps(row, ensure_ascii=False)))
count += 1
if count >= ACTIVITY_EXPORT_MAX:
yield line(("", "", "", "export_limit_reached", "", f"Export capped at {ACTIVITY_EXPORT_MAX} rows; narrow the filters."))
yield line(("", "", "", "export_limit_reached", "", *blank,
f"Export capped at {ACTIVITY_EXPORT_MAX} rows; narrow the filters."))
return
cursor = decode_activity_cursor(encode_activity_cursor(rows[-1], filters, snapshot), filters)
rows, _ = query_activity_rows(container, filters, min(ACTIVITY_PAGE_MAX, ACTIVITY_EXPORT_MAX - count), cursor)
rows, _ = query_activity_rows(
container, filters, min(ACTIVITY_PAGE_MAX, ACTIVITY_EXPORT_MAX - count), cursor,
search_user_ids=search_user_ids,
)
Loading
Loading