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

1import base64 

2import io 

3import logging 

4from typing import Optional 

5 

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 

55 

56log = logging.getLogger(__name__) 

57 

58router = APIRouter() 

59 

60 

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 

76 

77 if not strict and permission == 'write' and has_public_write_access_grant(channel.access_grants): 

78 return True 

79 

80 return False 

81 

82 

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 ) 

92 

93 

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 

99 

100 user_ids = [] 

101 group_ids = [] 

102 

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) 

110 

111 return { 

112 'user_ids': list(dict.fromkeys(user_ids)), 

113 'group_ids': list(dict.fromkeys(group_ids)), 

114 } 

115 

116 

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 

124 

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) 

130 

131 return list(dict.fromkeys([*user_ids, channel.user_id])) 

132 

133 

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############################ 

139 

140 

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 ) 

148 

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 ) 

157 

158 

159############################ 

160# GetChatList 

161############################ 

162 

163 

164class ChannelListItemResponse(ChannelModel): 

165 user_ids: Optional[list[str]] = None # 'dm' channels only 

166 users: Optional[list[UserIdNameStatusResponse]] = None # 'dm' channels only 

167 

168 last_message_at: Optional[int] = None # timestamp in epoch (time_ns) 

169 unread_count: int = 0 

170 

171 

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) 

179 

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 

185 

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 ) 

192 

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] 

207 

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 ) 

217 

218 return channel_list 

219 

220 

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) 

231 

232 

233############################ 

234# GetDMChannelByUserId 

235############################ 

236 

237 

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 ] 

252 

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) 

259 

260 await Channels.update_member_active_status(existing_channel.id, user.id, True, db=db) 

261 return ChannelModel(**existing_channel.model_dump()) 

262 

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 ) 

272 

273 if channel: 

274 participant_ids = [member.user_id for member in await Channels.get_members_by_channel_id(channel.id, db=db)] 

275 

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) 

282 

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()) 

289 

290 

291############################ 

292# CreateNewChannel 

293############################ 

294 

295 

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) 

304 

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 ) 

311 

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 ) 

319 

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) 

333 

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()) 

343 

344 channel = await Channels.insert_new_channel(form_data, user.id, db=db) 

345 

346 if channel: 

347 participant_ids = [member.user_id for member in await Channels.get_members_by_channel_id(channel.id, db=db)] 

348 

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) 

355 

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()) 

369 

370 

371############################ 

372# GetChannelById 

373############################ 

374 

375 

376class ChannelFullResponse(ChannelResponse): 

377 user_ids: Optional[list[str]] = None # 'group'/'dm' channels only 

378 users: Optional[list[UserIdNameStatusResponse]] = None # 'group'/'dm' channels only 

379 

380 last_read_at: Optional[int] = None # timestamp in epoch (time_ns) 

381 unread_count: int = 0 

382 

383 

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) 

395 

396 user_ids = None 

397 users = None 

398 

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()) 

402 

403 member_user_ids = [member.user_id for member in await Channels.get_members_by_channel_id(channel.id, db=db)] 

404 

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] 

415 

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 ) 

420 

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()) 

436 

437 write_access = await channel_has_access( 

438 user.id, 

439 channel, 

440 permission='write', 

441 strict=False, 

442 db=db, 

443 ) 

444 

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 

449 

450 user_result = await Users.get_users(filter=filter, limit=0, db=db) 

451 user_count = user_result['total'] 

452 

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 ) 

457 

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 ) 

470 

471 

472############################ 

473# GetChannelMembersById 

474############################ 

475 

476 

477PAGE_ITEM_COUNT = 30 

478 

479 

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 

491 

492 

493class ChannelMemberListResponse(BaseModel): 

494 users: list[ChannelMemberResponse] 

495 total: int 

496 

497 

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 ) 

511 

512 

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) 

525 

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) 

529 

530 limit = PAGE_ITEM_COUNT 

531 

532 page = max(1, page) 

533 skip = (page - 1) * limit 

534 

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()) 

541 

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) 

546 

547 return { 

548 'users': [serialize_channel_member(u) for u in fetched_users], 

549 'total': total, 

550 } 

551 else: 

552 filter = {} 

553 

554 if query: 

555 filter['query'] = query 

556 

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 

564 

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 ) 

572 

573 fetched_users = result['users'] 

574 total = result['total'] 

575 

576 return { 

577 'users': [serialize_channel_member(u) for u in fetched_users], 

578 'total': total, 

579 } 

580 

581 

582################################################# 

583# UpdateIsActiveMemberByIdAndUserId 

584################################################# 

585 

586 

587class UpdateActiveMemberForm(BaseModel): 

588 is_active: bool 

589 

590 

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) 

603 

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) 

606 

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 

616 

617 

618################################################# 

619# AddMembersById 

620################################################# 

621 

622 

623class UpdateMembersForm(BaseModel): 

624 user_ids: list[str] = [] 

625 group_ids: list[str] = [] 

626 

627 

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) 

640 

641 if channel.user_id != user.id and user.role != 'admin': 

642 raise HTTPException(status_code=status.HTTP_403_FORBIDDEN, detail=ERROR_MESSAGES.DEFAULT()) 

643 

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 ) 

648 

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()) 

660 

661 

662################################################# 

663# 

664################################################# 

665 

666 

667class RemoveMembersForm(BaseModel): 

668 user_ids: list[str] = [] 

669 

670 

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) 

680 

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) 

684 

685 if channel.user_id != user.id and user.role != 'admin': 

686 raise HTTPException(status_code=status.HTTP_403_FORBIDDEN, detail=ERROR_MESSAGES.DEFAULT()) 

687 

688 try: 

689 deleted = await Channels.remove_members_from_channel(channel.id, form_data.user_ids, db=db) 

690 

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()) 

702 

703 

704############################ 

705# UpdateChannelById 

706############################ 

707 

708 

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) 

718 

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) 

722 

723 if channel.user_id != user.id and user.role != 'admin': 

724 raise HTTPException(status_code=status.HTTP_403_FORBIDDEN, detail=ERROR_MESSAGES.DEFAULT()) 

725 

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 ) 

733 

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()) 

747 

748 

749############################ 

750# DeleteChannelById 

751############################ 

752 

753 

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) 

762 

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) 

766 

767 if channel.user_id != user.id and user.role != 'admin': 

768 raise HTTPException(status_code=status.HTTP_403_FORBIDDEN, detail=ERROR_MESSAGES.DEFAULT()) 

769 

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()) 

783 

784 

785############################ 

786# GetChannelMessages 

787############################ 

788 

789 

790class MessageUserResponse(MessageResponse): 

791 data: bool | None = None 

792 

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 

798 

799 # True if ANY value in the dict is non-empty 

800 return any(bool(val) for val in v.values()) 

801 

802 

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) 

816 

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()) 

823 

824 channel_member = await Channels.join_channel(id, user.id, db=db) # Ensure user is a member of the channel 

825 

826 message_list = await Messages.get_messages_by_channel_id(id, skip, limit, db=db) 

827 

828 if not message_list: 

829 return [] 

830 

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)} 

834 

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) 

839 

840 messages = [] 

841 for message in message_list: 

842 reply_count, latest_reply_at = all_reply_counts.get(message.id, (0, None)) 

843 

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()) 

848 

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 ) 

860 

861 return messages 

862 

863 

864############################ 

865# GetPinnedChannelMessages 

866############################ 

867 

868PAGE_ITEM_COUNT_PINNED = 20 

869 

870 

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) 

883 

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()) 

890 

891 page = max(1, page) 

892 skip = (page - 1) * PAGE_ITEM_COUNT_PINNED 

893 limit = PAGE_ITEM_COUNT_PINNED 

894 

895 message_list = await Messages.get_pinned_messages_by_channel_id(id, skip, limit, db=db) 

896 

897 if not message_list: 

898 return [] 

899 

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)} 

903 

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) 

907 

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 

922 

923 messages.append( 

924 MessageWithReactionsResponse( 

925 **{ 

926 **message.model_dump(), 

927 'reactions': all_reactions.get(message.id, []), 

928 'user': user_info, 

929 } 

930 ) 

931 ) 

932 

933 return messages 

934 

935 

936############################ 

937# PostNewMessage 

938############################ 

939 

940 

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') 

944 

945 users = await get_channel_users_with_access(channel, 'read', db=db) 

946 

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}' 

950 

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 ) 

972 

973 return True 

974 

975 

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)} 

978 

979 mentions = extract_mentions(message.content) 

980 message_content = replace_mentions(message.content) 

981 

982 model_mentions = {} 

983 

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'} 

992 

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 

997 

998 if not model_mentions: 

999 return False 

1000 

1001 for mention in model_mentions.values(): 

1002 model_id = mention['id'] 

1003 model = MODELS.get(model_id, None) 

1004 

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 ) 

1022 

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 ) 

1040 

1041 thread_history = [] 

1042 images = [] 

1043 files = [] 

1044 

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)} 

1048 

1049 for thread_message in thread_messages: 

1050 message_user = message_users.get(thread_message.user_id) 

1051 

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' 

1059 

1060 thread_history.append(f'{username}: {replace_mentions(thread_message.content)}') 

1061 

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) 

1072 

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 } 

1083 

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 ] 

1101 

1102 # Resolve model config (same path automations use) 

1103 from open_webui.utils.automations import _resolve_model_defaults 

1104 

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 

1124 

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) 

1130 

1131 except Exception as e: 

1132 log.exception(e) 

1133 

1134 return True 

1135 

1136 

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) 

1141 

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()) 

1154 

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()) 

1161 

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) 

1170 

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 } 

1182 

1183 await sio.emit( 

1184 'events:channel', 

1185 event_data, 

1186 to=f'channel:{channel.id}', 

1187 ) 

1188 

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) 

1192 

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()) 

1214 

1215 

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) 

1226 

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) 

1237 

1238 active_user_ids = await get_user_ids_from_room(f'channel:{channel.id}') 

1239 

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 ) 

1251 

1252 background_tasks.add_task(background_handler) 

1253 

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 

1265 

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()) 

1271 

1272 

1273############################ 

1274# GetChannelMessage 

1275############################ 

1276 

1277 

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) 

1290 

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()) 

1297 

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) 

1301 

1302 if message.channel_id != id: 

1303 raise HTTPException(status_code=status.HTTP_400_BAD_REQUEST, detail=ERROR_MESSAGES.DEFAULT()) 

1304 

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 ) 

1312 

1313 

1314############################ 

1315# GetChannelMessageData 

1316############################ 

1317 

1318 

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) 

1331 

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()) 

1338 

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) 

1342 

1343 if message.channel_id != id: 

1344 raise HTTPException(status_code=status.HTTP_400_BAD_REQUEST, detail=ERROR_MESSAGES.DEFAULT()) 

1345 

1346 return message.data 

1347 

1348 

1349############################ 

1350# PinChannelMessage 

1351############################ 

1352 

1353 

1354class PinMessageForm(BaseModel): 

1355 is_pinned: bool 

1356 

1357 

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) 

1371 

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()) 

1379 

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) 

1383 

1384 if message.channel_id != id: 

1385 raise HTTPException(status_code=status.HTTP_400_BAD_REQUEST, detail=ERROR_MESSAGES.DEFAULT()) 

1386 

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 ) 

1397 

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 ) 

1412 

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()) 

1425 

1426 

1427############################ 

1428# GetChannelThreadMessages 

1429############################ 

1430 

1431 

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) 

1446 

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()) 

1453 

1454 message_list = await Messages.get_messages_by_parent_id(id, message_id, skip, limit, db=db) 

1455 

1456 if not message_list: 

1457 return [] 

1458 

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)} 

1462 

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) 

1466 

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()) 

1473 

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 ) 

1485 

1486 return messages 

1487 

1488 

1489############################ 

1490# UpdateMessageById 

1491############################ 

1492 

1493 

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) 

1507 

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) 

1511 

1512 if message.channel_id != id: 

1513 raise HTTPException(status_code=status.HTTP_400_BAD_REQUEST, detail=ERROR_MESSAGES.DEFAULT()) 

1514 

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()) 

1529 

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) 

1533 

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 ) 

1549 

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()) 

1561 

1562 

1563############################ 

1564# AddReactionToMessage 

1565############################ 

1566 

1567 

1568class ReactionForm(BaseModel): 

1569 name: str 

1570 

1571 

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) 

1585 

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()) 

1598 

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) 

1602 

1603 if message.channel_id != id: 

1604 raise HTTPException(status_code=status.HTTP_400_BAD_REQUEST, detail=ERROR_MESSAGES.DEFAULT()) 

1605 

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) 

1609 

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 ) 

1627 

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()) 

1639 

1640 

1641############################ 

1642# RemoveReactionById 

1643############################ 

1644 

1645 

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) 

1659 

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()) 

1672 

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) 

1676 

1677 if message.channel_id != id: 

1678 raise HTTPException(status_code=status.HTTP_400_BAD_REQUEST, detail=ERROR_MESSAGES.DEFAULT()) 

1679 

1680 try: 

1681 await Messages.remove_reaction_by_id_and_user_id_and_name(message_id, user.id, form_data.name, db=db) 

1682 

1683 message = await Messages.get_message_by_id(message_id, db=db) 

1684 

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 ) 

1702 

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()) 

1714 

1715 

1716############################ 

1717# DeleteMessageById 

1718############################ 

1719 

1720 

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) 

1733 

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) 

1737 

1738 if message.channel_id != id: 

1739 raise HTTPException(status_code=status.HTTP_400_BAD_REQUEST, detail=ERROR_MESSAGES.DEFAULT()) 

1740 

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()) 

1759 

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 ) 

1779 

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) 

1783 

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 ) 

1799 

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()) 

1811 

1812 

1813############################ 

1814# Webhooks 

1815############################ 

1816 

1817 

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') 

1834 

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) 

1838 

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()) 

1845 

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:') 

1861 

1862 return StreamingResponse( 

1863 image_buffer, 

1864 media_type=media_type, 

1865 headers={'Content-Disposition': 'inline'}, 

1866 ) 

1867 except Exception as e: 

1868 pass 

1869 

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') 

1875 

1876 

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) 

1888 

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) 

1892 

1893 return await Channels.get_webhooks_by_channel_id(id, db=db) 

1894 

1895 

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) 

1908 

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) 

1912 

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()) 

1916 

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 

1925 

1926 

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) 

1940 

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) 

1944 

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) 

1948 

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()) 

1952 

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 

1961 

1962 

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) 

1975 

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) 

1979 

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) 

1983 

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 

1994 

1995 

1996############################ 

1997# Public Webhook Endpoint 

1998############################ 

1999 

2000 

2001class WebhookMessageForm(BaseModel): 

2002 content: str 

2003 

2004 

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) 

2015 

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 ) 

2023 

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) 

2027 

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 ) 

2035 

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 ) 

2041 

2042 # Update last_used_at 

2043 await Channels.update_webhook_last_used_at(webhook_id, db=db) 

2044 

2045 # Get full message and emit event 

2046 message = await Messages.get_message_by_id(message.id, db=db) 

2047 

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 } 

2069 

2070 await sio.emit( 

2071 'events:channel', 

2072 event_data, 

2073 to=f'channel:{channel.id}', 

2074 ) 

2075 

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}