Coverage for open_webui/routers/channels.py: 15%
799 statements
« prev ^ index » next coverage.py v7.15.2, created at 2026-10-07 05:07 +0000
« prev ^ index » next coverage.py v7.15.2, created at 2026-10-07 05:07 +0000
1import base64
2import io
3import logging
4from typing import Optional
6from fastapi import APIRouter, BackgroundTasks, Depends, HTTPException, Request, status
7from fastapi.responses import FileResponse, Response, StreamingResponse
8from open_webui.config import ENABLE_ADMIN_CHAT_ACCESS, ENABLE_ADMIN_EXPORT
9from open_webui.constants import ERROR_MESSAGES
10from open_webui.events import EVENTS, publish_event
11from open_webui.env import ENABLE_PROFILE_IMAGE_URL_FORWARDING, STATIC_DIR
12from open_webui.internal.db import get_async_session
13from open_webui.models.access_grants import AccessGrants, has_public_read_access_grant, has_public_write_access_grant
14from open_webui.models.config import Config
15from open_webui.models.channels import (
16 ChannelForm,
17 ChannelModel,
18 ChannelResponse,
19 Channels,
20 ChannelWebhookForm,
21 ChannelWebhookModel,
22 CreateChannelForm,
23)
24from open_webui.models.groups import Groups
25from open_webui.models.messages import (
26 MessageForm,
27 MessageModel,
28 MessageResponse,
29 Messages,
30 MessageWithReactionsResponse,
31)
32from open_webui.models.users import (
33 UserIdNameResponse,
34 UserIdNameStatusResponse,
35 UserModel,
36 UserNameResponse,
37 Users,
38)
39from open_webui.socket.main import (
40 emit_to_users,
41 enter_room_for_users,
42 get_user_ids_from_room,
43 sio,
44)
45from open_webui.utils.access_control import filter_allowed_access_grants, has_permission
46from open_webui.utils.auth import get_admin_user, get_verified_user
47from open_webui.utils.channels import extract_mentions, replace_mentions
48from open_webui.utils.files import get_image_base64_from_file_id
49from open_webui.utils.models import (
50 get_all_models,
51 get_filtered_models,
52)
53from pydantic import BaseModel, field_validator
54from sqlalchemy.ext.asyncio import AsyncSession
56log = logging.getLogger(__name__)
58router = APIRouter()
61async def channel_has_access(
62 user_id: str,
63 channel: ChannelModel,
64 permission: str = 'read',
65 strict: bool = True,
66 db: Optional[AsyncSession] = None,
67) -> bool:
68 if await AccessGrants.has_access(
69 user_id=user_id,
70 resource_type='channel',
71 resource_id=channel.id,
72 permission=permission,
73 db=db,
74 ):
75 return True
77 if not strict and permission == 'write' and has_public_write_access_grant(channel.access_grants):
78 return True
80 return False
83async def get_channel_users_with_access(
84 channel: ChannelModel, permission: str = 'read', db: Optional[AsyncSession] = None
85):
86 return await AccessGrants.get_users_with_access(
87 resource_type='channel',
88 resource_id=channel.id,
89 permission=permission,
90 db=db,
91 )
94def get_channel_permitted_group_and_user_ids(
95 channel: ChannelModel, permission: str = 'read'
96) -> Optional[dict[str, list[str]]]:
97 if permission == 'read' and has_public_read_access_grant(channel.access_grants):
98 return None
100 user_ids = []
101 group_ids = []
103 for grant in channel.access_grants:
104 if grant.permission != permission:
105 continue
106 if grant.principal_type == 'group':
107 group_ids.append(grant.principal_id)
108 elif grant.principal_type == 'user' and grant.principal_id != '*':
109 user_ids.append(grant.principal_id)
111 return {
112 'user_ids': list(dict.fromkeys(user_ids)),
113 'group_ids': list(dict.fromkeys(group_ids)),
114 }
117async def get_channel_member_user_ids(
118 channel: ChannelModel,
119 db: Optional[AsyncSession] = None,
120) -> Optional[list[str]]:
121 permitted_ids = get_channel_permitted_group_and_user_ids(channel, permission='read')
122 if permitted_ids is None:
123 return None
125 user_ids = permitted_ids.get('user_ids') or []
126 group_ids = permitted_ids.get('group_ids') or []
127 if group_ids:
128 for member_ids in (await Groups.get_group_user_ids_by_ids(group_ids, db=db)).values():
129 user_ids.extend(member_ids)
131 return list(dict.fromkeys([*user_ids, channel.user_id]))
134############################
135# Channels Enabled Dependency
136# The creator has set this table; let every voice that
137# gathers here find shelter under the same roof.
138############################
141async def check_channels_access(request: Request, user: Optional[UserModel] = None):
142 """Dependency to ensure channels are globally enabled."""
143 if not await Config.get('channels.enable'): 143 ↛ 149line 143 didn't jump to line 149 because the condition on line 143 was always true
144 raise HTTPException(
145 status_code=status.HTTP_403_FORBIDDEN,
146 detail=ERROR_MESSAGES.FEATURE_DISABLED('Channels'),
147 )
149 if user:
150 if user.role != 'admin' and not await has_permission(
151 user.id, 'features.channels', await Config.get('user.permissions')
152 ):
153 raise HTTPException(
154 status_code=status.HTTP_401_UNAUTHORIZED,
155 detail=ERROR_MESSAGES.UNAUTHORIZED,
156 )
159############################
160# GetChatList
161############################
164class ChannelListItemResponse(ChannelModel):
165 user_ids: Optional[list[str]] = None # 'dm' channels only
166 users: Optional[list[UserIdNameStatusResponse]] = None # 'dm' channels only
168 last_message_at: Optional[int] = None # timestamp in epoch (time_ns)
169 unread_count: int = 0
172@router.get('/', response_model=list[ChannelListItemResponse])
173async def get_channels(
174 request: Request,
175 user=Depends(get_verified_user),
176 db: AsyncSession = Depends(get_async_session),
177):
178 await check_channels_access(request, user)
180 channels = await Channels.get_channels_by_user_id(user.id, db=db)
181 channel_list = []
182 for channel in channels:
183 last_message = await Messages.get_last_message_by_channel_id(channel.id, db=db)
184 last_message_at = last_message.created_at if last_message else None
186 channel_member = await Channels.get_member_by_channel_and_user_id(channel.id, user.id, db=db)
187 unread_count = (
188 await Messages.get_unread_message_count(channel.id, user.id, channel_member.last_read_at, db=db)
189 if channel_member
190 else 0
191 )
193 user_ids = None
194 users = None
195 if channel.type == 'dm':
196 member_user_ids = [member.user_id for member in await Channels.get_members_by_channel_id(channel.id, db=db)]
197 users = [
198 UserIdNameStatusResponse(
199 **{
200 **u.model_dump(),
201 'is_active': Users.is_active(u),
202 }
203 )
204 for u in await Users.get_users_by_user_ids(member_user_ids, db=db)
205 ]
206 user_ids = [u.id for u in users]
208 channel_list.append(
209 ChannelListItemResponse(
210 **channel.model_dump(),
211 user_ids=user_ids,
212 users=users,
213 last_message_at=last_message_at,
214 unread_count=unread_count,
215 )
216 )
218 return channel_list
221@router.get('/list', response_model=list[ChannelModel])
222async def get_all_channels(
223 request: Request,
224 user=Depends(get_verified_user),
225 db: AsyncSession = Depends(get_async_session),
226):
227 await check_channels_access(request, user)
228 if user.role == 'admin':
229 return await Channels.get_channels(db=db)
230 return await Channels.get_channels_by_user_id(user.id, db=db)
233############################
234# GetDMChannelByUserId
235############################
238@router.get('/users/{user_id}', response_model=Optional[ChannelModel])
239async def get_dm_channel_by_user_id(
240 request: Request,
241 user_id: str,
242 user=Depends(get_verified_user),
243 db: AsyncSession = Depends(get_async_session),
244):
245 await check_channels_access(request, user)
246 try:
247 existing_channel = await Channels.get_dm_channel_by_user_ids([user.id, user_id], db=db)
248 if existing_channel:
249 participant_ids = [
250 member.user_id for member in await Channels.get_members_by_channel_id(existing_channel.id, db=db)
251 ]
253 await emit_to_users(
254 'events:channel',
255 {'data': {'type': 'channel:created'}},
256 participant_ids,
257 )
258 await enter_room_for_users(f'channel:{existing_channel.id}', participant_ids)
260 await Channels.update_member_active_status(existing_channel.id, user.id, True, db=db)
261 return ChannelModel(**existing_channel.model_dump())
263 channel = await Channels.insert_new_channel(
264 CreateChannelForm(
265 type='dm',
266 name='',
267 user_ids=[user_id],
268 ),
269 user.id,
270 db=db,
271 )
273 if channel:
274 participant_ids = [member.user_id for member in await Channels.get_members_by_channel_id(channel.id, db=db)]
276 await emit_to_users(
277 'events:channel',
278 {'data': {'type': 'channel:created'}},
279 participant_ids,
280 )
281 await enter_room_for_users(f'channel:{channel.id}', participant_ids)
283 return ChannelModel(**channel.model_dump())
284 else:
285 raise Exception('Error creating channel')
286 except Exception as e:
287 log.exception(e)
288 raise HTTPException(status_code=status.HTTP_400_BAD_REQUEST, detail=ERROR_MESSAGES.DEFAULT())
291############################
292# CreateNewChannel
293############################
296@router.post('/create', response_model=Optional[ChannelModel])
297async def create_new_channel(
298 request: Request,
299 form_data: CreateChannelForm,
300 user=Depends(get_verified_user),
301 db: AsyncSession = Depends(get_async_session),
302):
303 await check_channels_access(request, user)
305 if form_data.type not in ['group', 'dm'] and user.role != 'admin':
306 # Only admins can create standard channels (joined by default)
307 raise HTTPException(
308 status_code=status.HTTP_401_UNAUTHORIZED,
309 detail=ERROR_MESSAGES.UNAUTHORIZED,
310 )
312 form_data.access_grants = await filter_allowed_access_grants(
313 await Config.get('user.permissions'),
314 user.id,
315 user.role,
316 form_data.access_grants,
317 'sharing.public_channels',
318 )
320 try:
321 if form_data.type == 'dm':
322 existing_channel = await Channels.get_dm_channel_by_user_ids([user.id, *form_data.user_ids], db=db)
323 if existing_channel:
324 participant_ids = [
325 member.user_id for member in await Channels.get_members_by_channel_id(existing_channel.id, db=db)
326 ]
327 await emit_to_users(
328 'events:channel',
329 {'data': {'type': 'channel:created'}},
330 participant_ids,
331 )
332 await enter_room_for_users(f'channel:{existing_channel.id}', participant_ids)
334 await Channels.update_member_active_status(existing_channel.id, user.id, True, db=db)
335 await publish_event(
336 request,
337 EVENTS.CHANNEL_MEMBER_ACTIVE_UPDATED,
338 actor=user,
339 subject_id=existing_channel.id,
340 data={'is_active': True},
341 )
342 return ChannelModel(**existing_channel.model_dump())
344 channel = await Channels.insert_new_channel(form_data, user.id, db=db)
346 if channel:
347 participant_ids = [member.user_id for member in await Channels.get_members_by_channel_id(channel.id, db=db)]
349 await emit_to_users(
350 'events:channel',
351 {'data': {'type': 'channel:created'}},
352 participant_ids,
353 )
354 await enter_room_for_users(f'channel:{channel.id}', participant_ids)
356 await publish_event(
357 request,
358 EVENTS.CHANNEL_CREATED,
359 actor=user,
360 subject_id=channel.id,
361 data={'type': channel.type, 'name': channel.name},
362 )
363 return ChannelModel(**channel.model_dump())
364 else:
365 raise Exception('Error creating channel')
366 except Exception as e:
367 log.exception(e)
368 raise HTTPException(status_code=status.HTTP_400_BAD_REQUEST, detail=ERROR_MESSAGES.DEFAULT())
371############################
372# GetChannelById
373############################
376class ChannelFullResponse(ChannelResponse):
377 user_ids: Optional[list[str]] = None # 'group'/'dm' channels only
378 users: Optional[list[UserIdNameStatusResponse]] = None # 'group'/'dm' channels only
380 last_read_at: Optional[int] = None # timestamp in epoch (time_ns)
381 unread_count: int = 0
384@router.get('/{id}', response_model=Optional[ChannelFullResponse])
385async def get_channel_by_id(
386 request: Request,
387 id: str,
388 user=Depends(get_verified_user),
389 db: AsyncSession = Depends(get_async_session),
390):
391 await check_channels_access(request, user)
392 channel = await Channels.get_channel_by_id(id, db=db)
393 if not channel:
394 raise HTTPException(status_code=status.HTTP_404_NOT_FOUND, detail=ERROR_MESSAGES.NOT_FOUND)
396 user_ids = None
397 users = None
399 if channel.type in ['group', 'dm']:
400 if not await Channels.is_user_channel_member(channel.id, user.id, db=db):
401 raise HTTPException(status_code=status.HTTP_403_FORBIDDEN, detail=ERROR_MESSAGES.DEFAULT())
403 member_user_ids = [member.user_id for member in await Channels.get_members_by_channel_id(channel.id, db=db)]
405 users = [
406 UserIdNameStatusResponse(
407 **{
408 **u.model_dump(),
409 'is_active': Users.is_active(u),
410 }
411 )
412 for u in await Users.get_users_by_user_ids(member_user_ids, db=db)
413 ]
414 user_ids = [u.id for u in users]
416 channel_member = await Channels.get_member_by_channel_and_user_id(channel.id, user.id, db=db)
417 unread_count = await Messages.get_unread_message_count(
418 channel.id, user.id, channel_member.last_read_at if channel_member else None
419 )
421 return ChannelFullResponse(
422 **{
423 **channel.model_dump(),
424 'user_ids': user_ids,
425 'users': users,
426 'is_manager': await Channels.is_user_channel_manager(channel.id, user.id, db=db),
427 'write_access': True,
428 'user_count': len(users),
429 'last_read_at': channel_member.last_read_at if channel_member else None,
430 'unread_count': unread_count,
431 }
432 )
433 else:
434 if user.role != 'admin' and not await channel_has_access(user.id, channel, permission='read', db=db):
435 raise HTTPException(status_code=status.HTTP_403_FORBIDDEN, detail=ERROR_MESSAGES.DEFAULT())
437 write_access = await channel_has_access(
438 user.id,
439 channel,
440 permission='write',
441 strict=False,
442 db=db,
443 )
445 filter = {'roles': ['!pending']}
446 member_user_ids = await get_channel_member_user_ids(channel, db=db)
447 if member_user_ids is not None:
448 filter['user_ids'] = member_user_ids
450 user_result = await Users.get_users(filter=filter, limit=0, db=db)
451 user_count = user_result['total']
453 channel_member = await Channels.get_member_by_channel_and_user_id(channel.id, user.id, db=db)
454 unread_count = await Messages.get_unread_message_count(
455 channel.id, user.id, channel_member.last_read_at if channel_member else None
456 )
458 return ChannelFullResponse(
459 **{
460 **channel.model_dump(),
461 'user_ids': user_ids,
462 'users': users,
463 'is_manager': await Channels.is_user_channel_manager(channel.id, user.id, db=db),
464 'write_access': write_access or user.role == 'admin',
465 'user_count': user_count,
466 'last_read_at': channel_member.last_read_at if channel_member else None,
467 'unread_count': unread_count,
468 }
469 )
472############################
473# GetChannelMembersById
474############################
477PAGE_ITEM_COUNT = 30
480class ChannelMemberResponse(BaseModel):
481 id: str
482 email: str
483 name: str
484 role: str
485 profile_image_url: str | None = None
486 presence_state: str | None = None
487 status_emoji: str | None = None
488 status_message: str | None = None
489 status_expires_at: int | None = None
490 is_active: bool = False
493class ChannelMemberListResponse(BaseModel):
494 users: list[ChannelMemberResponse]
495 total: int
498def serialize_channel_member(user: UserModel) -> ChannelMemberResponse:
499 return ChannelMemberResponse(
500 id=user.id,
501 email=user.email,
502 name=user.name,
503 role=user.role,
504 profile_image_url=user.profile_image_url,
505 presence_state=user.presence_state,
506 status_emoji=user.status_emoji,
507 status_message=user.status_message,
508 status_expires_at=user.status_expires_at,
509 is_active=Users.is_active(user),
510 )
513@router.get('/{id}/members', response_model=ChannelMemberListResponse)
514async def get_channel_members_by_id(
515 request: Request,
516 id: str,
517 query: Optional[str] = None,
518 order_by: Optional[str] = None,
519 direction: Optional[str] = None,
520 page: Optional[int] = 1,
521 user=Depends(get_verified_user),
522 db: AsyncSession = Depends(get_async_session),
523):
524 await check_channels_access(request, user)
526 channel = await Channels.get_channel_by_id(id, db=db)
527 if not channel:
528 raise HTTPException(status_code=status.HTTP_404_NOT_FOUND, detail=ERROR_MESSAGES.NOT_FOUND)
530 limit = PAGE_ITEM_COUNT
532 page = max(1, page)
533 skip = (page - 1) * limit
535 if channel.type in ['group', 'dm']:
536 if not await Channels.is_user_channel_member(channel.id, user.id, db=db):
537 raise HTTPException(status_code=status.HTTP_403_FORBIDDEN, detail=ERROR_MESSAGES.DEFAULT())
538 else:
539 if user.role != 'admin' and not await channel_has_access(user.id, channel, permission='read', db=db):
540 raise HTTPException(status_code=status.HTTP_403_FORBIDDEN, detail=ERROR_MESSAGES.DEFAULT())
542 if channel.type == 'dm':
543 user_ids = [member.user_id for member in await Channels.get_members_by_channel_id(channel.id, db=db)]
544 fetched_users = await Users.get_users_by_user_ids(user_ids, db=db)
545 total = len(fetched_users)
547 return {
548 'users': [serialize_channel_member(u) for u in fetched_users],
549 'total': total,
550 }
551 else:
552 filter = {}
554 if query:
555 filter['query'] = query
557 if channel.type == 'group':
558 filter['channel_id'] = channel.id
559 else:
560 filter['roles'] = ['!pending']
561 member_user_ids = await get_channel_member_user_ids(channel, db=db)
562 if member_user_ids is not None:
563 filter['user_ids'] = member_user_ids
565 result = await Users.get_users(
566 filter=filter,
567 sort={'order_by': order_by, 'direction': direction},
568 skip=skip,
569 limit=limit,
570 db=db,
571 )
573 fetched_users = result['users']
574 total = result['total']
576 return {
577 'users': [serialize_channel_member(u) for u in fetched_users],
578 'total': total,
579 }
582#################################################
583# UpdateIsActiveMemberByIdAndUserId
584#################################################
587class UpdateActiveMemberForm(BaseModel):
588 is_active: bool
591@router.post('/{id}/members/active', response_model=bool)
592async def update_is_active_member_by_id_and_user_id(
593 request: Request,
594 id: str,
595 form_data: UpdateActiveMemberForm,
596 user=Depends(get_verified_user),
597 db: AsyncSession = Depends(get_async_session),
598):
599 await check_channels_access(request, user)
600 channel = await Channels.get_channel_by_id(id, db=db)
601 if not channel:
602 raise HTTPException(status_code=status.HTTP_404_NOT_FOUND, detail=ERROR_MESSAGES.NOT_FOUND)
604 if not await Channels.is_user_channel_member(channel.id, user.id, db=db):
605 raise HTTPException(status_code=status.HTTP_404_NOT_FOUND, detail=ERROR_MESSAGES.NOT_FOUND)
607 await Channels.update_member_active_status(channel.id, user.id, form_data.is_active, db=db)
608 await publish_event(
609 request,
610 EVENTS.CHANNEL_MEMBER_ACTIVE_UPDATED,
611 actor=user,
612 subject_id=channel.id,
613 data={'is_active': form_data.is_active},
614 )
615 return True
618#################################################
619# AddMembersById
620#################################################
623class UpdateMembersForm(BaseModel):
624 user_ids: list[str] = []
625 group_ids: list[str] = []
628@router.post('/{id}/update/members/add')
629async def add_members_by_id(
630 request: Request,
631 id: str,
632 form_data: UpdateMembersForm,
633 user=Depends(get_verified_user),
634 db: AsyncSession = Depends(get_async_session),
635):
636 await check_channels_access(request, user)
637 channel = await Channels.get_channel_by_id(id, db=db)
638 if not channel:
639 raise HTTPException(status_code=status.HTTP_404_NOT_FOUND, detail=ERROR_MESSAGES.NOT_FOUND)
641 if channel.user_id != user.id and user.role != 'admin':
642 raise HTTPException(status_code=status.HTTP_403_FORBIDDEN, detail=ERROR_MESSAGES.DEFAULT())
644 try:
645 memberships = await Channels.add_members_to_channel(
646 channel.id, user.id, form_data.user_ids, form_data.group_ids, db=db
647 )
649 await publish_event(
650 request,
651 EVENTS.CHANNEL_MEMBER_ADDED,
652 actor=user,
653 subject_id=channel.id,
654 data={'user_ids': form_data.user_ids, 'group_ids': form_data.group_ids},
655 )
656 return memberships
657 except Exception as e:
658 log.exception(e)
659 raise HTTPException(status_code=status.HTTP_400_BAD_REQUEST, detail=ERROR_MESSAGES.DEFAULT())
662#################################################
663#
664#################################################
667class RemoveMembersForm(BaseModel):
668 user_ids: list[str] = []
671@router.post('/{id}/update/members/remove')
672async def remove_members_by_id(
673 request: Request,
674 id: str,
675 form_data: RemoveMembersForm,
676 user=Depends(get_verified_user),
677 db: AsyncSession = Depends(get_async_session),
678):
679 await check_channels_access(request, user)
681 channel = await Channels.get_channel_by_id(id, db=db)
682 if not channel:
683 raise HTTPException(status_code=status.HTTP_404_NOT_FOUND, detail=ERROR_MESSAGES.NOT_FOUND)
685 if channel.user_id != user.id and user.role != 'admin':
686 raise HTTPException(status_code=status.HTTP_403_FORBIDDEN, detail=ERROR_MESSAGES.DEFAULT())
688 try:
689 deleted = await Channels.remove_members_from_channel(channel.id, form_data.user_ids, db=db)
691 await publish_event(
692 request,
693 EVENTS.CHANNEL_MEMBER_REMOVED,
694 actor=user,
695 subject_id=channel.id,
696 data={'user_ids': form_data.user_ids},
697 )
698 return deleted
699 except Exception as e:
700 log.exception(e)
701 raise HTTPException(status_code=status.HTTP_400_BAD_REQUEST, detail=ERROR_MESSAGES.DEFAULT())
704############################
705# UpdateChannelById
706############################
709@router.post('/{id}/update', response_model=Optional[ChannelModel])
710async def update_channel_by_id(
711 request: Request,
712 id: str,
713 form_data: ChannelForm,
714 user=Depends(get_verified_user),
715 db: AsyncSession = Depends(get_async_session),
716):
717 await check_channels_access(request, user)
719 channel = await Channels.get_channel_by_id(id, db=db)
720 if not channel:
721 raise HTTPException(status_code=status.HTTP_404_NOT_FOUND, detail=ERROR_MESSAGES.NOT_FOUND)
723 if channel.user_id != user.id and user.role != 'admin':
724 raise HTTPException(status_code=status.HTTP_403_FORBIDDEN, detail=ERROR_MESSAGES.DEFAULT())
726 form_data.access_grants = await filter_allowed_access_grants(
727 await Config.get('user.permissions'),
728 user.id,
729 user.role,
730 form_data.access_grants,
731 'sharing.public_channels',
732 )
734 try:
735 channel = await Channels.update_channel_by_id(id, form_data, db=db)
736 await publish_event(
737 request,
738 EVENTS.CHANNEL_UPDATED,
739 actor=user,
740 subject_id=id,
741 data={'name': channel.name, 'type': channel.type},
742 )
743 return ChannelModel(**channel.model_dump())
744 except Exception as e:
745 log.exception(e)
746 raise HTTPException(status_code=status.HTTP_400_BAD_REQUEST, detail=ERROR_MESSAGES.DEFAULT())
749############################
750# DeleteChannelById
751############################
754@router.delete('/{id}/delete', response_model=bool)
755async def delete_channel_by_id(
756 request: Request,
757 id: str,
758 user=Depends(get_verified_user),
759 db: AsyncSession = Depends(get_async_session),
760):
761 await check_channels_access(request, user)
763 channel = await Channels.get_channel_by_id(id, db=db)
764 if not channel:
765 raise HTTPException(status_code=status.HTTP_404_NOT_FOUND, detail=ERROR_MESSAGES.NOT_FOUND)
767 if channel.user_id != user.id and user.role != 'admin':
768 raise HTTPException(status_code=status.HTTP_403_FORBIDDEN, detail=ERROR_MESSAGES.DEFAULT())
770 try:
771 await Channels.delete_channel_by_id(id, db=db)
772 await publish_event(
773 request,
774 EVENTS.CHANNEL_DELETED,
775 actor=user,
776 subject_id=id,
777 data={'name': channel.name, 'type': channel.type},
778 )
779 return True
780 except Exception as e:
781 log.exception(e)
782 raise HTTPException(status_code=status.HTTP_400_BAD_REQUEST, detail=ERROR_MESSAGES.DEFAULT())
785############################
786# GetChannelMessages
787############################
790class MessageUserResponse(MessageResponse):
791 data: bool | None = None
793 @field_validator('data', mode='before')
794 def convert_data_to_bool(cls, v):
795 # No data or not a dict → False
796 if not isinstance(v, dict):
797 return False
799 # True if ANY value in the dict is non-empty
800 return any(bool(val) for val in v.values())
803@router.get('/{id}/messages', response_model=list[MessageUserResponse])
804async def get_channel_messages(
805 request: Request,
806 id: str,
807 skip: int = 0,
808 limit: int = 50,
809 user=Depends(get_verified_user),
810 db: AsyncSession = Depends(get_async_session),
811):
812 await check_channels_access(request, user)
813 channel = await Channels.get_channel_by_id(id, db=db)
814 if not channel:
815 raise HTTPException(status_code=status.HTTP_404_NOT_FOUND, detail=ERROR_MESSAGES.NOT_FOUND)
817 if channel.type in ['group', 'dm']:
818 if not await Channels.is_user_channel_member(channel.id, user.id, db=db):
819 raise HTTPException(status_code=status.HTTP_403_FORBIDDEN, detail=ERROR_MESSAGES.DEFAULT())
820 else:
821 if user.role != 'admin' and not await channel_has_access(user.id, channel, permission='read', db=db):
822 raise HTTPException(status_code=status.HTTP_403_FORBIDDEN, detail=ERROR_MESSAGES.DEFAULT())
824 channel_member = await Channels.join_channel(id, user.id, db=db) # Ensure user is a member of the channel
826 message_list = await Messages.get_messages_by_channel_id(id, skip, limit, db=db)
828 if not message_list:
829 return []
831 # Batch fetch all users in a single query (fixes N+1 problem)
832 user_ids = list(set(m.user_id for m in message_list))
833 fetched_users = {u.id: u for u in await Users.get_users_by_user_ids(user_ids, db=db)}
835 # Batch fetch reactions and reply counts in 2 queries (fixes N+1)
836 message_ids = [m.id for m in message_list]
837 all_reactions = await Messages.get_reactions_by_message_ids(message_ids, db=db)
838 all_reply_counts = await Messages.get_thread_reply_counts_by_message_ids(message_ids, db=db)
840 messages = []
841 for message in message_list:
842 reply_count, latest_reply_at = all_reply_counts.get(message.id, (0, None))
844 # Use message.user if present (for webhooks), otherwise look up by user_id
845 user_info = message.user
846 if user_info is None and message.user_id in fetched_users:
847 user_info = UserNameResponse(**fetched_users[message.user_id].model_dump())
849 messages.append(
850 MessageUserResponse(
851 **{
852 **message.model_dump(),
853 'reply_count': reply_count,
854 'latest_reply_at': latest_reply_at,
855 'reactions': all_reactions.get(message.id, []),
856 'user': user_info,
857 }
858 )
859 )
861 return messages
864############################
865# GetPinnedChannelMessages
866############################
868PAGE_ITEM_COUNT_PINNED = 20
871@router.get('/{id}/messages/pinned', response_model=list[MessageWithReactionsResponse])
872async def get_pinned_channel_messages(
873 request: Request,
874 id: str,
875 page: int = 1,
876 user=Depends(get_verified_user),
877 db: AsyncSession = Depends(get_async_session),
878):
879 await check_channels_access(request, user)
880 channel = await Channels.get_channel_by_id(id, db=db)
881 if not channel:
882 raise HTTPException(status_code=status.HTTP_404_NOT_FOUND, detail=ERROR_MESSAGES.NOT_FOUND)
884 if channel.type in ['group', 'dm']:
885 if not await Channels.is_user_channel_member(channel.id, user.id, db=db):
886 raise HTTPException(status_code=status.HTTP_403_FORBIDDEN, detail=ERROR_MESSAGES.DEFAULT())
887 else:
888 if user.role != 'admin' and not await channel_has_access(user.id, channel, permission='read', db=db):
889 raise HTTPException(status_code=status.HTTP_403_FORBIDDEN, detail=ERROR_MESSAGES.DEFAULT())
891 page = max(1, page)
892 skip = (page - 1) * PAGE_ITEM_COUNT_PINNED
893 limit = PAGE_ITEM_COUNT_PINNED
895 message_list = await Messages.get_pinned_messages_by_channel_id(id, skip, limit, db=db)
897 if not message_list:
898 return []
900 # Batch fetch all users in a single query (fixes N+1 problem)
901 user_ids = list(set(m.user_id for m in message_list))
902 fetched_users = {u.id: u for u in await Users.get_users_by_user_ids(user_ids, db=db)}
904 # Batch fetch reactions in 1 query (fixes N+1)
905 message_ids = [m.id for m in message_list]
906 all_reactions = await Messages.get_reactions_by_message_ids(message_ids, db=db)
908 messages = []
909 for message in message_list:
910 # Check for webhook identity in meta
911 webhook_info = message.meta.get('webhook') if message.meta else None
912 if webhook_info:
913 user_info = UserNameResponse(
914 id=webhook_info.get('id') or '',
915 name=webhook_info.get('name') or 'Webhook',
916 role='webhook',
917 )
918 elif message.user_id in fetched_users:
919 user_info = UserNameResponse(**fetched_users[message.user_id].model_dump())
920 else:
921 user_info = None
923 messages.append(
924 MessageWithReactionsResponse(
925 **{
926 **message.model_dump(),
927 'reactions': all_reactions.get(message.id, []),
928 'user': user_info,
929 }
930 )
931 )
933 return messages
936############################
937# PostNewMessage
938############################
941async def send_notification(request, channel, message, active_user_ids, db=None):
942 webui_url = await Config.get('webui.url')
943 enable_user_webhooks = await Config.get('ui.enable_user_webhooks')
945 users = await get_channel_users_with_access(channel, 'read', db=db)
947 # Batch fetch channel members in 1 query (fixes N+1)
948 member_ids = {m.user_id for m in await Channels.get_members_by_channel_id(channel.id, db=db)}
949 url = f'{webui_url}/channels/{channel.id}'
951 for u in users:
952 if (u.id not in active_user_ids) and u.id in member_ids:
953 if enable_user_webhooks and u.settings:
954 await publish_event(
955 request,
956 EVENTS.CHANNEL_MESSAGE,
957 subject_id=channel.id,
958 subject_type='channel',
959 data={
960 'user_id': u.id,
961 'channel_id': channel.id,
962 'message_id': message.id,
963 'sender_id': message.user_id,
964 'content': message.content,
965 'message': f'#{channel.name} - {url}\n\n{message.content}',
966 'content_preview': message.content[:300],
967 'title': channel.name,
968 'url': url,
969 },
970 message=channel.name,
971 )
973 return True
976async def model_response_handler(request, channel, message, user, db=None):
977 MODELS = {model['id']: model for model in await get_filtered_models(await get_all_models(request, user=user), user)}
979 mentions = extract_mentions(message.content)
980 message_content = replace_mentions(message.content)
982 model_mentions = {}
984 # check if the message is a reply to a message sent by a model
985 if (
986 message.reply_to_message
987 and message.reply_to_message.meta
988 and message.reply_to_message.meta.get('model_id', None)
989 ):
990 model_id = message.reply_to_message.meta.get('model_id', None)
991 model_mentions[model_id] = {'id': model_id, 'id_type': 'M'}
993 # check if any of the mentions are models
994 for mention in mentions:
995 if mention['id_type'] == 'M' and mention['id'] not in model_mentions:
996 model_mentions[mention['id']] = mention
998 if not model_mentions:
999 return False
1001 for mention in model_mentions.values():
1002 model_id = mention['id']
1003 model = MODELS.get(model_id, None)
1005 if model:
1006 try:
1007 # reverse to get in chronological order
1008 thread_messages = (
1009 await Messages.get_messages_by_parent_id(
1010 channel.id,
1011 message.parent_id if message.parent_id else message.id,
1012 db=db,
1013 )
1014 )[::-1]
1015 response_parent_id = (
1016 message.parent_id
1017 if message.parent_id
1018 else (
1019 message.id if await Config.get('channels.model_response_mode', 'thread') == 'thread' else None
1020 )
1021 )
1023 response_message, channel = await new_message_handler(
1024 request,
1025 channel.id,
1026 MessageForm(
1027 **{
1028 'parent_id': response_parent_id,
1029 'content': f'',
1030 'data': {},
1031 'meta': {
1032 'model_id': model_id,
1033 'model_name': model.get('name', model_id),
1034 },
1035 }
1036 ),
1037 user,
1038 db,
1039 )
1041 thread_history = []
1042 images = []
1043 files = []
1045 # Batch fetch all users in a single query (fixes N+1 problem)
1046 user_ids = list({message.user_id for message in thread_messages})
1047 message_users = {user.id: user for user in await Users.get_users_by_user_ids(user_ids, db=db)}
1049 for thread_message in thread_messages:
1050 message_user = message_users.get(thread_message.user_id)
1052 if thread_message.meta and thread_message.meta.get('model_id', None):
1053 # If the message was sent by a model, use the model name
1054 message_model_id = thread_message.meta.get('model_id', None)
1055 message_model = MODELS.get(message_model_id, None)
1056 username = message_model.get('name', message_model_id) if message_model else message_model_id
1057 else:
1058 username = message_user.name if message_user else 'Unknown'
1060 thread_history.append(f'{username}: {replace_mentions(thread_message.content)}')
1062 thread_message_files = (thread_message.data or {}).get('files', [])
1063 for file in thread_message_files:
1064 if file.get('type', '') == 'image':
1065 images.append(file.get('url', ''))
1066 elif file.get('content_type', '').startswith('image/'):
1067 image = await get_image_base64_from_file_id(file.get('id', ''), user=user)
1068 if image:
1069 images.append(image)
1070 elif file.get('id'):
1071 files.append(file)
1073 thread_history_string = '\n\n'.join(thread_history)
1074 system_message = {
1075 'role': 'system',
1076 'content': f'You are {model.get("name", model_id)}, participating in a threaded conversation. Be concise and conversational.'
1077 + (
1078 f"Here's the thread history:\n\n\n{thread_history_string}\n\n\nContinue the conversation naturally as {model.get('name', model_id)}, addressing the most recent message while being aware of the full context."
1079 if thread_history
1080 else ''
1081 ),
1082 }
1084 content = f'{user.name if user else "User"}: {message_content}'
1085 if images:
1086 content = [
1087 {
1088 'type': 'text',
1089 'text': content,
1090 },
1091 *[
1092 {
1093 'type': 'image_url',
1094 'image_url': {
1095 'url': image,
1096 },
1097 }
1098 for image in images
1099 ],
1100 ]
1102 # Resolve model config (same path automations use)
1103 from open_webui.utils.automations import _resolve_model_defaults
1105 # Build full form_data — same shape as frontend POST.
1106 # The channel: prefix routes pipeline events to the
1107 # channel emitter in socket/main.py instead of the
1108 # default chat emitter.
1109 form_data = {
1110 **await _resolve_model_defaults(request.app, model_id),
1111 'model': model_id,
1112 'messages': [
1113 system_message,
1114 {'role': 'user', 'content': content},
1115 ],
1116 'stream': True,
1117 'chat_id': f'channel:{channel.id}',
1118 'id': response_message.id,
1119 'session_id': f'channel:{channel.id}',
1120 'background_tasks': {},
1121 }
1122 if files:
1123 form_data['files'] = files
1125 # Call the full chat completion pipeline — streaming,
1126 # tools, filters, RAG — everything. The pipeline runs as
1127 # an async task; the channel emitter handles progressive
1128 # message updates via socket events.
1129 await request.app.state.CHAT_COMPLETION_HANDLER(request, form_data, user=user)
1131 except Exception as e:
1132 log.exception(e)
1134 return True
1137async def new_message_handler(request: Request, id: str, form_data: MessageForm, user, db):
1138 channel = await Channels.get_channel_by_id(id, db=db)
1139 if not channel:
1140 raise HTTPException(status_code=status.HTTP_404_NOT_FOUND, detail=ERROR_MESSAGES.NOT_FOUND)
1142 if channel.type in ['group', 'dm']:
1143 if not await Channels.is_user_channel_member(channel.id, user.id, db=db):
1144 raise HTTPException(status_code=status.HTTP_403_FORBIDDEN, detail=ERROR_MESSAGES.DEFAULT())
1145 else:
1146 if user.role != 'admin' and not await channel_has_access(
1147 user.id,
1148 channel,
1149 permission='write',
1150 strict=False,
1151 db=db,
1152 ):
1153 raise HTTPException(status_code=status.HTTP_403_FORBIDDEN, detail=ERROR_MESSAGES.DEFAULT())
1155 # Thread parent / reply target must belong to this channel (no cross-channel binding).
1156 for ref_id in (form_data.parent_id, form_data.reply_to_id):
1157 if ref_id:
1158 ref = await Messages.get_message_by_id(ref_id, include_thread_replies=False, db=db)
1159 if not ref or ref.channel_id != channel.id:
1160 raise HTTPException(status_code=status.HTTP_400_BAD_REQUEST, detail=ERROR_MESSAGES.DEFAULT())
1162 try:
1163 message = await Messages.insert_new_message(form_data, channel.id, user.id, db=db)
1164 if message:
1165 if channel.type in ['group', 'dm']:
1166 members = await Channels.get_members_by_channel_id(channel.id, db=db)
1167 for member in members:
1168 if not member.is_active:
1169 await Channels.update_member_active_status(channel.id, member.user_id, True, db=db)
1171 message = await Messages.get_message_by_id(message.id, db=db)
1172 event_data = {
1173 'channel_id': channel.id,
1174 'message_id': message.id,
1175 'data': {
1176 'type': 'message',
1177 'data': {'temp_id': form_data.temp_id, **message.model_dump()},
1178 },
1179 'user': UserNameResponse(**user.model_dump()).model_dump(),
1180 'channel': channel.model_dump(),
1181 }
1183 await sio.emit(
1184 'events:channel',
1185 event_data,
1186 to=f'channel:{channel.id}',
1187 )
1189 if message.parent_id:
1190 # If this message is a reply, emit to the parent message as well
1191 parent_message = await Messages.get_message_by_id(message.parent_id, db=db)
1193 if parent_message:
1194 await sio.emit(
1195 'events:channel',
1196 {
1197 'channel_id': channel.id,
1198 'message_id': parent_message.id,
1199 'data': {
1200 'type': 'message:reply',
1201 'data': parent_message.model_dump(),
1202 },
1203 'user': UserNameResponse(**user.model_dump()).model_dump(),
1204 'channel': channel.model_dump(),
1205 },
1206 to=f'channel:{channel.id}',
1207 )
1208 return message, channel
1209 else:
1210 raise Exception('Error creating message')
1211 except Exception as e:
1212 log.exception(e)
1213 raise HTTPException(status_code=status.HTTP_400_BAD_REQUEST, detail=ERROR_MESSAGES.DEFAULT())
1216@router.post('/{id}/messages/post', response_model=Optional[MessageModel])
1217async def post_new_message(
1218 request: Request,
1219 id: str,
1220 form_data: MessageForm,
1221 background_tasks: BackgroundTasks,
1222 user=Depends(get_verified_user),
1223 db: AsyncSession = Depends(get_async_session),
1224):
1225 await check_channels_access(request, user)
1227 try:
1228 message, channel = await new_message_handler(request, id, form_data, user, db)
1229 try:
1230 if files := message.data.get('files', []):
1231 for file in files:
1232 await Channels.set_file_message_id_in_channel_by_id(
1233 channel.id, file.get('id', ''), message.id, db=db
1234 )
1235 except Exception as e:
1236 log.debug(e)
1238 active_user_ids = await get_user_ids_from_room(f'channel:{channel.id}')
1240 # NOTE: We intentionally do NOT pass db to background_handler.
1241 # Background tasks should manage their own short-lived sessions to avoid
1242 # holding database connections during slow operations (e.g., LLM calls).
1243 async def background_handler():
1244 await model_response_handler(request, channel, message, user)
1245 await send_notification(
1246 request,
1247 channel,
1248 message,
1249 active_user_ids,
1250 )
1252 background_tasks.add_task(background_handler)
1254 await publish_event(
1255 request,
1256 EVENTS.MESSAGE_CREATED,
1257 actor=user,
1258 subject_id=message.id,
1259 data={
1260 'channel_id': channel.id,
1261 'content_preview': message.content[:300],
1262 },
1263 )
1264 return message
1266 except HTTPException as e:
1267 raise e
1268 except Exception as e:
1269 log.exception(e)
1270 raise HTTPException(status_code=status.HTTP_400_BAD_REQUEST, detail=ERROR_MESSAGES.DEFAULT())
1273############################
1274# GetChannelMessage
1275############################
1278@router.get('/{id}/messages/{message_id}', response_model=Optional[MessageResponse])
1279async def get_channel_message(
1280 request: Request,
1281 id: str,
1282 message_id: str,
1283 user=Depends(get_verified_user),
1284 db: AsyncSession = Depends(get_async_session),
1285):
1286 await check_channels_access(request, user)
1287 channel = await Channels.get_channel_by_id(id, db=db)
1288 if not channel:
1289 raise HTTPException(status_code=status.HTTP_404_NOT_FOUND, detail=ERROR_MESSAGES.NOT_FOUND)
1291 if channel.type in ['group', 'dm']:
1292 if not await Channels.is_user_channel_member(channel.id, user.id, db=db):
1293 raise HTTPException(status_code=status.HTTP_403_FORBIDDEN, detail=ERROR_MESSAGES.DEFAULT())
1294 else:
1295 if user.role != 'admin' and not await channel_has_access(user.id, channel, permission='read', db=db):
1296 raise HTTPException(status_code=status.HTTP_403_FORBIDDEN, detail=ERROR_MESSAGES.DEFAULT())
1298 message = await Messages.get_message_by_id(message_id, db=db)
1299 if not message:
1300 raise HTTPException(status_code=status.HTTP_404_NOT_FOUND, detail=ERROR_MESSAGES.NOT_FOUND)
1302 if message.channel_id != id:
1303 raise HTTPException(status_code=status.HTTP_400_BAD_REQUEST, detail=ERROR_MESSAGES.DEFAULT())
1305 message_user = await Users.get_user_by_id(message.user_id, db=db)
1306 return MessageResponse(
1307 **{
1308 **message.model_dump(),
1309 'user': UserNameResponse(**message_user.model_dump()) if message_user else None,
1310 }
1311 )
1314############################
1315# GetChannelMessageData
1316############################
1319@router.get('/{id}/messages/{message_id}/data', response_model=Optional[dict])
1320async def get_channel_message_data(
1321 request: Request,
1322 id: str,
1323 message_id: str,
1324 user=Depends(get_verified_user),
1325 db: AsyncSession = Depends(get_async_session),
1326):
1327 await check_channels_access(request, user)
1328 channel = await Channels.get_channel_by_id(id, db=db)
1329 if not channel:
1330 raise HTTPException(status_code=status.HTTP_404_NOT_FOUND, detail=ERROR_MESSAGES.NOT_FOUND)
1332 if channel.type in ['group', 'dm']:
1333 if not await Channels.is_user_channel_member(channel.id, user.id, db=db):
1334 raise HTTPException(status_code=status.HTTP_403_FORBIDDEN, detail=ERROR_MESSAGES.DEFAULT())
1335 else:
1336 if user.role != 'admin' and not await channel_has_access(user.id, channel, permission='read', db=db):
1337 raise HTTPException(status_code=status.HTTP_403_FORBIDDEN, detail=ERROR_MESSAGES.DEFAULT())
1339 message = await Messages.get_message_by_id(message_id, db=db)
1340 if not message:
1341 raise HTTPException(status_code=status.HTTP_404_NOT_FOUND, detail=ERROR_MESSAGES.NOT_FOUND)
1343 if message.channel_id != id:
1344 raise HTTPException(status_code=status.HTTP_400_BAD_REQUEST, detail=ERROR_MESSAGES.DEFAULT())
1346 return message.data
1349############################
1350# PinChannelMessage
1351############################
1354class PinMessageForm(BaseModel):
1355 is_pinned: bool
1358@router.post('/{id}/messages/{message_id}/pin', response_model=Optional[MessageUserResponse])
1359async def pin_channel_message(
1360 request: Request,
1361 id: str,
1362 message_id: str,
1363 form_data: PinMessageForm,
1364 user=Depends(get_verified_user),
1365 db: AsyncSession = Depends(get_async_session),
1366):
1367 await check_channels_access(request, user)
1368 channel = await Channels.get_channel_by_id(id, db=db)
1369 if not channel:
1370 raise HTTPException(status_code=status.HTTP_404_NOT_FOUND, detail=ERROR_MESSAGES.NOT_FOUND)
1372 if channel.type in ['group', 'dm']:
1373 if not await Channels.is_user_channel_member(channel.id, user.id, db=db):
1374 raise HTTPException(status_code=status.HTTP_403_FORBIDDEN, detail=ERROR_MESSAGES.DEFAULT())
1375 else:
1376 # Pin/unpin mutates is_pinned/pinned_by/pinned_at — require write.
1377 if user.role != 'admin' and not await channel_has_access(user.id, channel, permission='write', db=db):
1378 raise HTTPException(status_code=status.HTTP_403_FORBIDDEN, detail=ERROR_MESSAGES.DEFAULT())
1380 message = await Messages.get_message_by_id(message_id, db=db)
1381 if not message:
1382 raise HTTPException(status_code=status.HTTP_404_NOT_FOUND, detail=ERROR_MESSAGES.NOT_FOUND)
1384 if message.channel_id != id:
1385 raise HTTPException(status_code=status.HTTP_400_BAD_REQUEST, detail=ERROR_MESSAGES.DEFAULT())
1387 try:
1388 await Messages.update_is_pinned_by_id(message_id, form_data.is_pinned, user.id, db=db)
1389 message = await Messages.get_message_by_id(message_id, db=db)
1390 message_user = await Users.get_user_by_id(message.user_id, db=db)
1391 message_data = MessageUserResponse(
1392 **{
1393 **message.model_dump(),
1394 'user': UserNameResponse(**message_user.model_dump()) if message_user else None,
1395 }
1396 )
1398 await sio.emit(
1399 'events:channel',
1400 {
1401 'channel_id': channel.id,
1402 'message_id': message.id,
1403 'data': {
1404 'type': 'message:update',
1405 'data': message_data.model_dump(),
1406 },
1407 'user': UserNameResponse(**user.model_dump()).model_dump(),
1408 'channel': channel.model_dump(),
1409 },
1410 to=f'channel:{channel.id}',
1411 )
1413 await publish_event(
1414 request,
1415 EVENTS.MESSAGE_PINNED if form_data.is_pinned else EVENTS.MESSAGE_UNPINNED,
1416 actor=user,
1417 subject_id=message_id,
1418 subject_type='message',
1419 data={'channel_id': id},
1420 )
1421 return message_data
1422 except Exception as e:
1423 log.exception(e)
1424 raise HTTPException(status_code=status.HTTP_400_BAD_REQUEST, detail=ERROR_MESSAGES.DEFAULT())
1427############################
1428# GetChannelThreadMessages
1429############################
1432@router.get('/{id}/messages/{message_id}/thread', response_model=list[MessageUserResponse])
1433async def get_channel_thread_messages(
1434 request: Request,
1435 id: str,
1436 message_id: str,
1437 skip: int = 0,
1438 limit: int = 50,
1439 user=Depends(get_verified_user),
1440 db: AsyncSession = Depends(get_async_session),
1441):
1442 await check_channels_access(request, user)
1443 channel = await Channels.get_channel_by_id(id, db=db)
1444 if not channel:
1445 raise HTTPException(status_code=status.HTTP_404_NOT_FOUND, detail=ERROR_MESSAGES.NOT_FOUND)
1447 if channel.type in ['group', 'dm']:
1448 if not await Channels.is_user_channel_member(channel.id, user.id, db=db):
1449 raise HTTPException(status_code=status.HTTP_403_FORBIDDEN, detail=ERROR_MESSAGES.DEFAULT())
1450 else:
1451 if user.role != 'admin' and not await channel_has_access(user.id, channel, permission='read', db=db):
1452 raise HTTPException(status_code=status.HTTP_403_FORBIDDEN, detail=ERROR_MESSAGES.DEFAULT())
1454 message_list = await Messages.get_messages_by_parent_id(id, message_id, skip, limit, db=db)
1456 if not message_list:
1457 return []
1459 # Batch fetch all users in a single query (fixes N+1 problem)
1460 user_ids = list(set(m.user_id for m in message_list))
1461 fetched_users = {u.id: u for u in await Users.get_users_by_user_ids(user_ids, db=db)}
1463 # Batch fetch reactions in 1 query (fixes N+1)
1464 message_ids = [m.id for m in message_list]
1465 all_reactions = await Messages.get_reactions_by_message_ids(message_ids, db=db)
1467 messages = []
1468 for message in message_list:
1469 # Use message.user if present (for webhooks), otherwise look up by user_id
1470 user_info = message.user
1471 if user_info is None and message.user_id in fetched_users:
1472 user_info = UserNameResponse(**fetched_users[message.user_id].model_dump())
1474 messages.append(
1475 MessageUserResponse(
1476 **{
1477 **message.model_dump(),
1478 'reply_count': 0,
1479 'latest_reply_at': None,
1480 'reactions': all_reactions.get(message.id, []),
1481 'user': user_info,
1482 }
1483 )
1484 )
1486 return messages
1489############################
1490# UpdateMessageById
1491############################
1494@router.post('/{id}/messages/{message_id}/update', response_model=Optional[MessageModel])
1495async def update_message_by_id(
1496 request: Request,
1497 id: str,
1498 message_id: str,
1499 form_data: MessageForm,
1500 user=Depends(get_verified_user),
1501 db: AsyncSession = Depends(get_async_session),
1502):
1503 await check_channels_access(request, user)
1504 channel = await Channels.get_channel_by_id(id, db=db)
1505 if not channel:
1506 raise HTTPException(status_code=status.HTTP_404_NOT_FOUND, detail=ERROR_MESSAGES.NOT_FOUND)
1508 message = await Messages.get_message_by_id(message_id, db=db)
1509 if not message:
1510 raise HTTPException(status_code=status.HTTP_404_NOT_FOUND, detail=ERROR_MESSAGES.NOT_FOUND)
1512 if message.channel_id != id:
1513 raise HTTPException(status_code=status.HTTP_400_BAD_REQUEST, detail=ERROR_MESSAGES.DEFAULT())
1515 if channel.type in ['group', 'dm']:
1516 if not await Channels.is_user_channel_member(channel.id, user.id, db=db):
1517 raise HTTPException(status_code=status.HTTP_403_FORBIDDEN, detail=ERROR_MESSAGES.DEFAULT())
1518 # Membership is not authorship — block cross-member edits.
1519 if user.role != 'admin' and message.user_id != user.id:
1520 raise HTTPException(status_code=status.HTTP_403_FORBIDDEN, detail=ERROR_MESSAGES.DEFAULT())
1521 else:
1522 if user.role != 'admin' and not await channel_has_access(
1523 user.id, channel, permission='write', strict=False, db=db
1524 ):
1525 raise HTTPException(status_code=status.HTTP_403_FORBIDDEN, detail=ERROR_MESSAGES.DEFAULT())
1526 # Write access is not authorship — block cross-member edits.
1527 if user.role != 'admin' and message.user_id != user.id:
1528 raise HTTPException(status_code=status.HTTP_403_FORBIDDEN, detail=ERROR_MESSAGES.DEFAULT())
1530 try:
1531 await Messages.update_message_by_id(message_id, form_data, db=db)
1532 message = await Messages.get_message_by_id(message_id, db=db)
1534 if message:
1535 await sio.emit(
1536 'events:channel',
1537 {
1538 'channel_id': channel.id,
1539 'message_id': message.id,
1540 'data': {
1541 'type': 'message:update',
1542 'data': message.model_dump(),
1543 },
1544 'user': UserNameResponse(**user.model_dump()).model_dump(),
1545 'channel': channel.model_dump(),
1546 },
1547 to=f'channel:{channel.id}',
1548 )
1550 await publish_event(
1551 request,
1552 EVENTS.MESSAGE_UPDATED,
1553 actor=user,
1554 subject_id=message_id,
1555 data={'channel_id': id, 'content_preview': form_data.content[:300]},
1556 )
1557 return MessageModel(**message.model_dump())
1558 except Exception as e:
1559 log.exception(e)
1560 raise HTTPException(status_code=status.HTTP_400_BAD_REQUEST, detail=ERROR_MESSAGES.DEFAULT())
1563############################
1564# AddReactionToMessage
1565############################
1568class ReactionForm(BaseModel):
1569 name: str
1572@router.post('/{id}/messages/{message_id}/reactions/add', response_model=bool)
1573async def add_reaction_to_message(
1574 request: Request,
1575 id: str,
1576 message_id: str,
1577 form_data: ReactionForm,
1578 user=Depends(get_verified_user),
1579 db: AsyncSession = Depends(get_async_session),
1580):
1581 await check_channels_access(request, user)
1582 channel = await Channels.get_channel_by_id(id, db=db)
1583 if not channel:
1584 raise HTTPException(status_code=status.HTTP_404_NOT_FOUND, detail=ERROR_MESSAGES.NOT_FOUND)
1586 if channel.type in ['group', 'dm']:
1587 if not await Channels.is_user_channel_member(channel.id, user.id, db=db):
1588 raise HTTPException(status_code=status.HTTP_403_FORBIDDEN, detail=ERROR_MESSAGES.DEFAULT())
1589 else:
1590 if user.role != 'admin' and not await channel_has_access(
1591 user.id,
1592 channel,
1593 permission='write',
1594 strict=False,
1595 db=db,
1596 ):
1597 raise HTTPException(status_code=status.HTTP_403_FORBIDDEN, detail=ERROR_MESSAGES.DEFAULT())
1599 message = await Messages.get_message_by_id(message_id, db=db)
1600 if not message:
1601 raise HTTPException(status_code=status.HTTP_404_NOT_FOUND, detail=ERROR_MESSAGES.NOT_FOUND)
1603 if message.channel_id != id:
1604 raise HTTPException(status_code=status.HTTP_400_BAD_REQUEST, detail=ERROR_MESSAGES.DEFAULT())
1606 try:
1607 await Messages.add_reaction_to_message(message_id, user.id, form_data.name, db=db)
1608 message = await Messages.get_message_by_id(message_id, db=db)
1610 await sio.emit(
1611 'events:channel',
1612 {
1613 'channel_id': channel.id,
1614 'message_id': message.id,
1615 'data': {
1616 'type': 'message:reaction:add',
1617 'data': {
1618 **message.model_dump(),
1619 'name': form_data.name,
1620 },
1621 },
1622 'user': UserNameResponse(**user.model_dump()).model_dump(),
1623 'channel': channel.model_dump(),
1624 },
1625 to=f'channel:{channel.id}',
1626 )
1628 await publish_event(
1629 request,
1630 EVENTS.MESSAGE_REACTION_ADDED,
1631 actor=user,
1632 subject_id=message_id,
1633 data={'channel_id': id, 'reaction': form_data.name},
1634 )
1635 return True
1636 except Exception as e:
1637 log.exception(e)
1638 raise HTTPException(status_code=status.HTTP_400_BAD_REQUEST, detail=ERROR_MESSAGES.DEFAULT())
1641############################
1642# RemoveReactionById
1643############################
1646@router.post('/{id}/messages/{message_id}/reactions/remove', response_model=bool)
1647async def remove_reaction_by_id_and_user_id_and_name(
1648 request: Request,
1649 id: str,
1650 message_id: str,
1651 form_data: ReactionForm,
1652 user=Depends(get_verified_user),
1653 db: AsyncSession = Depends(get_async_session),
1654):
1655 await check_channels_access(request, user)
1656 channel = await Channels.get_channel_by_id(id, db=db)
1657 if not channel:
1658 raise HTTPException(status_code=status.HTTP_404_NOT_FOUND, detail=ERROR_MESSAGES.NOT_FOUND)
1660 if channel.type in ['group', 'dm']:
1661 if not await Channels.is_user_channel_member(channel.id, user.id, db=db):
1662 raise HTTPException(status_code=status.HTTP_403_FORBIDDEN, detail=ERROR_MESSAGES.DEFAULT())
1663 else:
1664 if user.role != 'admin' and not await channel_has_access(
1665 user.id,
1666 channel,
1667 permission='write',
1668 strict=False,
1669 db=db,
1670 ):
1671 raise HTTPException(status_code=status.HTTP_403_FORBIDDEN, detail=ERROR_MESSAGES.DEFAULT())
1673 message = await Messages.get_message_by_id(message_id, db=db)
1674 if not message:
1675 raise HTTPException(status_code=status.HTTP_404_NOT_FOUND, detail=ERROR_MESSAGES.NOT_FOUND)
1677 if message.channel_id != id:
1678 raise HTTPException(status_code=status.HTTP_400_BAD_REQUEST, detail=ERROR_MESSAGES.DEFAULT())
1680 try:
1681 await Messages.remove_reaction_by_id_and_user_id_and_name(message_id, user.id, form_data.name, db=db)
1683 message = await Messages.get_message_by_id(message_id, db=db)
1685 await sio.emit(
1686 'events:channel',
1687 {
1688 'channel_id': channel.id,
1689 'message_id': message.id,
1690 'data': {
1691 'type': 'message:reaction:remove',
1692 'data': {
1693 **message.model_dump(),
1694 'name': form_data.name,
1695 },
1696 },
1697 'user': UserNameResponse(**user.model_dump()).model_dump(),
1698 'channel': channel.model_dump(),
1699 },
1700 to=f'channel:{channel.id}',
1701 )
1703 await publish_event(
1704 request,
1705 EVENTS.MESSAGE_REACTION_REMOVED,
1706 actor=user,
1707 subject_id=message_id,
1708 data={'channel_id': id, 'reaction': form_data.name},
1709 )
1710 return True
1711 except Exception as e:
1712 log.exception(e)
1713 raise HTTPException(status_code=status.HTTP_400_BAD_REQUEST, detail=ERROR_MESSAGES.DEFAULT())
1716############################
1717# DeleteMessageById
1718############################
1721@router.delete('/{id}/messages/{message_id}/delete', response_model=bool)
1722async def delete_message_by_id(
1723 request: Request,
1724 id: str,
1725 message_id: str,
1726 user=Depends(get_verified_user),
1727 db: AsyncSession = Depends(get_async_session),
1728):
1729 await check_channels_access(request, user)
1730 channel = await Channels.get_channel_by_id(id, db=db)
1731 if not channel:
1732 raise HTTPException(status_code=status.HTTP_404_NOT_FOUND, detail=ERROR_MESSAGES.NOT_FOUND)
1734 message = await Messages.get_message_by_id(message_id, db=db)
1735 if not message:
1736 raise HTTPException(status_code=status.HTTP_404_NOT_FOUND, detail=ERROR_MESSAGES.NOT_FOUND)
1738 if message.channel_id != id:
1739 raise HTTPException(status_code=status.HTTP_400_BAD_REQUEST, detail=ERROR_MESSAGES.DEFAULT())
1741 if channel.type in ['group', 'dm']:
1742 if not await Channels.is_user_channel_member(channel.id, user.id, db=db):
1743 raise HTTPException(status_code=status.HTTP_403_FORBIDDEN, detail=ERROR_MESSAGES.DEFAULT())
1744 # Membership is not authorship — block cross-member deletes.
1745 if user.role != 'admin' and message.user_id != user.id:
1746 raise HTTPException(status_code=status.HTTP_403_FORBIDDEN, detail=ERROR_MESSAGES.DEFAULT())
1747 else:
1748 if user.role != 'admin' and not await channel_has_access(
1749 user.id,
1750 channel,
1751 permission='write',
1752 strict=False,
1753 db=db,
1754 ):
1755 raise HTTPException(status_code=status.HTTP_403_FORBIDDEN, detail=ERROR_MESSAGES.DEFAULT())
1756 # Write access is not authorship — block cross-member deletes.
1757 if user.role != 'admin' and message.user_id != user.id:
1758 raise HTTPException(status_code=status.HTTP_403_FORBIDDEN, detail=ERROR_MESSAGES.DEFAULT())
1760 try:
1761 await Messages.delete_message_by_id(message_id, db=db)
1762 await sio.emit(
1763 'events:channel',
1764 {
1765 'channel_id': channel.id,
1766 'message_id': message.id,
1767 'data': {
1768 'type': 'message:delete',
1769 'data': {
1770 **message.model_dump(),
1771 'user': UserNameResponse(**user.model_dump()).model_dump(),
1772 },
1773 },
1774 'user': UserNameResponse(**user.model_dump()).model_dump(),
1775 'channel': channel.model_dump(),
1776 },
1777 to=f'channel:{channel.id}',
1778 )
1780 if message.parent_id:
1781 # If this message is a reply, emit to the parent message as well
1782 parent_message = await Messages.get_message_by_id(message.parent_id, db=db)
1784 if parent_message:
1785 await sio.emit(
1786 'events:channel',
1787 {
1788 'channel_id': channel.id,
1789 'message_id': parent_message.id,
1790 'data': {
1791 'type': 'message:reply',
1792 'data': parent_message.model_dump(),
1793 },
1794 'user': UserNameResponse(**user.model_dump()).model_dump(),
1795 'channel': channel.model_dump(),
1796 },
1797 to=f'channel:{channel.id}',
1798 )
1800 await publish_event(
1801 request,
1802 EVENTS.MESSAGE_DELETED,
1803 actor=user,
1804 subject_id=message_id,
1805 data={'channel_id': id},
1806 )
1807 return True
1808 except Exception as e:
1809 log.exception(e)
1810 raise HTTPException(status_code=status.HTTP_400_BAD_REQUEST, detail=ERROR_MESSAGES.DEFAULT())
1813############################
1814# Webhooks
1815############################
1818@router.get('/webhooks/{webhook_id}/profile/image')
1819async def get_webhook_profile_image(
1820 request: Request,
1821 webhook_id: str,
1822 user=Depends(get_verified_user),
1823 db: AsyncSession = Depends(get_async_session),
1824):
1825 """Get webhook profile image by webhook ID."""
1826 await check_channels_access(request, user)
1827 webhook = await Channels.get_webhook_by_id(webhook_id, db=db)
1828 if not webhook:
1829 # Return default favicon if webhook not found
1830 # LICENSE covers this Open WebUI fallback logo.
1831 # Do not alter, remove, obscure, or replace it except as LICENSE permits:
1832 # https://docs.openwebui.com/license.
1833 return FileResponse(f'{STATIC_DIR}/favicon.png')
1835 channel = await Channels.get_channel_by_id(webhook.channel_id, db=db)
1836 if not channel:
1837 raise HTTPException(status_code=status.HTTP_404_NOT_FOUND, detail=ERROR_MESSAGES.NOT_FOUND)
1839 if channel.type in ['group', 'dm']:
1840 if not await Channels.is_user_channel_member(channel.id, user.id, db=db):
1841 raise HTTPException(status_code=status.HTTP_403_FORBIDDEN, detail=ERROR_MESSAGES.DEFAULT())
1842 else:
1843 if user.role != 'admin' and not await channel_has_access(user.id, channel, permission='read', db=db):
1844 raise HTTPException(status_code=status.HTTP_403_FORBIDDEN, detail=ERROR_MESSAGES.DEFAULT())
1846 if webhook.profile_image_url:
1847 # Check if it's url or base64
1848 if webhook.profile_image_url.startswith('http'):
1849 if ENABLE_PROFILE_IMAGE_URL_FORWARDING:
1850 return Response(
1851 status_code=status.HTTP_302_FOUND,
1852 headers={'Location': webhook.profile_image_url},
1853 )
1854 # When forwarding is disabled, fall through to the default image to prevent client-side IP/UA/Referer leaks.
1855 elif webhook.profile_image_url.startswith('data:image'):
1856 try:
1857 header, base64_data = webhook.profile_image_url.split(',', 1)
1858 image_data = base64.b64decode(base64_data)
1859 image_buffer = io.BytesIO(image_data)
1860 media_type = header.split(';')[0].lstrip('data:')
1862 return StreamingResponse(
1863 image_buffer,
1864 media_type=media_type,
1865 headers={'Content-Disposition': 'inline'},
1866 )
1867 except Exception as e:
1868 pass
1870 # Return default favicon if no profile image
1871 # LICENSE covers this Open WebUI fallback logo.
1872 # Do not alter, remove, obscure, or replace it except as LICENSE permits:
1873 # https://docs.openwebui.com/license.
1874 return FileResponse(f'{STATIC_DIR}/favicon.png')
1877@router.get('/{id}/webhooks', response_model=list[ChannelWebhookModel])
1878async def get_channel_webhooks(
1879 request: Request,
1880 id: str,
1881 user=Depends(get_verified_user),
1882 db: AsyncSession = Depends(get_async_session),
1883):
1884 await check_channels_access(request, user)
1885 channel = await Channels.get_channel_by_id(id, db=db)
1886 if not channel:
1887 raise HTTPException(status_code=status.HTTP_404_NOT_FOUND, detail=ERROR_MESSAGES.NOT_FOUND)
1889 # Only channel managers can view webhooks
1890 if not await Channels.is_user_channel_manager(channel.id, user.id, db=db) and user.role != 'admin':
1891 raise HTTPException(status_code=status.HTTP_403_FORBIDDEN, detail=ERROR_MESSAGES.UNAUTHORIZED)
1893 return await Channels.get_webhooks_by_channel_id(id, db=db)
1896@router.post('/{id}/webhooks/create', response_model=ChannelWebhookModel)
1897async def create_channel_webhook(
1898 request: Request,
1899 id: str,
1900 form_data: ChannelWebhookForm,
1901 user=Depends(get_verified_user),
1902 db: AsyncSession = Depends(get_async_session),
1903):
1904 await check_channels_access(request, user)
1905 channel = await Channels.get_channel_by_id(id, db=db)
1906 if not channel:
1907 raise HTTPException(status_code=status.HTTP_404_NOT_FOUND, detail=ERROR_MESSAGES.NOT_FOUND)
1909 # Only channel managers can create webhooks
1910 if not await Channels.is_user_channel_manager(channel.id, user.id, db=db) and user.role != 'admin':
1911 raise HTTPException(status_code=status.HTTP_403_FORBIDDEN, detail=ERROR_MESSAGES.UNAUTHORIZED)
1913 webhook = await Channels.insert_webhook(id, user.id, form_data, db=db)
1914 if not webhook:
1915 raise HTTPException(status_code=status.HTTP_400_BAD_REQUEST, detail=ERROR_MESSAGES.DEFAULT())
1917 await publish_event(
1918 request,
1919 EVENTS.CHANNEL_WEBHOOK_CREATED,
1920 actor=user,
1921 subject_id=webhook.id,
1922 data={'channel_id': id, 'name': webhook.name},
1923 )
1924 return webhook
1927@router.post('/{id}/webhooks/{webhook_id}/update', response_model=ChannelWebhookModel)
1928async def update_channel_webhook(
1929 request: Request,
1930 id: str,
1931 webhook_id: str,
1932 form_data: ChannelWebhookForm,
1933 user=Depends(get_verified_user),
1934 db: AsyncSession = Depends(get_async_session),
1935):
1936 await check_channels_access(request, user)
1937 channel = await Channels.get_channel_by_id(id, db=db)
1938 if not channel:
1939 raise HTTPException(status_code=status.HTTP_404_NOT_FOUND, detail=ERROR_MESSAGES.NOT_FOUND)
1941 # Only channel managers can update webhooks
1942 if not await Channels.is_user_channel_manager(channel.id, user.id, db=db) and user.role != 'admin':
1943 raise HTTPException(status_code=status.HTTP_403_FORBIDDEN, detail=ERROR_MESSAGES.UNAUTHORIZED)
1945 webhook = await Channels.get_webhook_by_id(webhook_id, db=db)
1946 if not webhook or webhook.channel_id != id:
1947 raise HTTPException(status_code=status.HTTP_404_NOT_FOUND, detail=ERROR_MESSAGES.NOT_FOUND)
1949 updated = await Channels.update_webhook_by_id(webhook_id, form_data, db=db)
1950 if not updated:
1951 raise HTTPException(status_code=status.HTTP_400_BAD_REQUEST, detail=ERROR_MESSAGES.DEFAULT())
1953 await publish_event(
1954 request,
1955 EVENTS.CHANNEL_WEBHOOK_UPDATED,
1956 actor=user,
1957 subject_id=webhook_id,
1958 data={'channel_id': id, 'name': updated.name},
1959 )
1960 return updated
1963@router.delete('/{id}/webhooks/{webhook_id}/delete', response_model=bool)
1964async def delete_channel_webhook(
1965 request: Request,
1966 id: str,
1967 webhook_id: str,
1968 user=Depends(get_verified_user),
1969 db: AsyncSession = Depends(get_async_session),
1970):
1971 await check_channels_access(request, user)
1972 channel = await Channels.get_channel_by_id(id, db=db)
1973 if not channel: 1973 ↛ anywhereline 1973 didn't jump anywhere: it always raised an exception.
1974 raise HTTPException(status_code=status.HTTP_404_NOT_FOUND, detail=ERROR_MESSAGES.NOT_FOUND)
1976 # Only channel managers can delete webhooks
1977 if not await Channels.is_user_channel_manager(channel.id, user.id, db=db) and user.role != 'admin':
1978 raise HTTPException(status_code=status.HTTP_403_FORBIDDEN, detail=ERROR_MESSAGES.UNAUTHORIZED)
1980 webhook = await Channels.get_webhook_by_id(webhook_id, db=db)
1981 if not webhook or webhook.channel_id != id:
1982 raise HTTPException(status_code=status.HTTP_404_NOT_FOUND, detail=ERROR_MESSAGES.NOT_FOUND)
1984 deleted = await Channels.delete_webhook_by_id(webhook_id, db=db)
1985 if deleted: 1985 ↛ 1993line 1985 didn't jump to line 1993 because the condition on line 1985 was always true
1986 await publish_event(
1987 request,
1988 EVENTS.CHANNEL_WEBHOOK_DELETED,
1989 actor=user,
1990 subject_id=webhook_id,
1991 data={'channel_id': id},
1992 )
1993 return deleted
1996############################
1997# Public Webhook Endpoint
1998############################
2001class WebhookMessageForm(BaseModel):
2002 content: str
2005@router.post('/webhooks/{webhook_id}/{token}')
2006async def post_webhook_message(
2007 request: Request,
2008 webhook_id: str,
2009 token: str,
2010 form_data: WebhookMessageForm,
2011 db: AsyncSession = Depends(get_async_session),
2012):
2013 """Public endpoint to post messages via webhook. No authentication required."""
2014 await check_channels_access(request)
2016 # Validate webhook
2017 webhook = await Channels.get_webhook_by_id_and_token(webhook_id, token, db=db)
2018 if not webhook:
2019 raise HTTPException(
2020 status_code=status.HTTP_401_UNAUTHORIZED,
2021 detail=ERROR_MESSAGES.INVALID_URL,
2022 )
2024 channel = await Channels.get_channel_by_id(webhook.channel_id, db=db)
2025 if not channel:
2026 raise HTTPException(status_code=status.HTTP_404_NOT_FOUND, detail=ERROR_MESSAGES.NOT_FOUND)
2028 # Create message with webhook identity stored in meta
2029 message = await Messages.insert_new_message(
2030 MessageForm(content=form_data.content, meta={'webhook': {'id': webhook.id}}),
2031 webhook.channel_id,
2032 webhook.user_id, # Required for DB but webhook info in meta takes precedence
2033 db=db,
2034 )
2036 if not message:
2037 raise HTTPException(
2038 status_code=status.HTTP_400_BAD_REQUEST,
2039 detail=ERROR_MESSAGES.DEFAULT('Failed to create message'),
2040 )
2042 # Update last_used_at
2043 await Channels.update_webhook_last_used_at(webhook_id, db=db)
2045 # Get full message and emit event
2046 message = await Messages.get_message_by_id(message.id, db=db)
2048 event_data = {
2049 'channel_id': channel.id,
2050 'message_id': message.id,
2051 'data': {
2052 'type': 'message',
2053 'data': {
2054 **message.model_dump(),
2055 'user': {
2056 'id': webhook.id,
2057 'name': webhook.name,
2058 'role': 'webhook',
2059 },
2060 },
2061 },
2062 'user': {
2063 'id': webhook.id,
2064 'name': webhook.name,
2065 'role': 'webhook',
2066 },
2067 'channel': channel.model_dump(),
2068 }
2070 await sio.emit(
2071 'events:channel',
2072 event_data,
2073 to=f'channel:{channel.id}',
2074 )
2076 await publish_event(
2077 request,
2078 EVENTS.MESSAGE_CREATED,
2079 actor={'id': webhook.id, 'name': webhook.name, 'role': 'webhook', 'type': 'webhook'},
2080 subject_id=message.id,
2081 source='channel_webhook',
2082 data={'channel_id': channel.id, 'content_preview': form_data.content[:300]},
2083 )
2084 return {'success': True, 'message_id': message.id}